FlinkSQL支持HANA数据库离线采集
生产中遇到一些譬如SAP系统使用HANA数据库,要求使用FlinkSQL进行离线采集。
记录一下源码修改:
1. 增加hana.dialect
HanaDialect.java
package org.apache.flink.connector.jdbc.databases.hana.dialect;
import org.apache.flink.connector.jdbc.converter.JdbcRowConverter;
import org.apache.flink.connector.jdbc.dialect.AbstractDialect;
import org.apache.flink.table.types.logical.LogicalTypeRoot;
import org.apache.flink.table.types.logical.RowType;
import java.util.EnumSet;
import java.util.Optional;
import java.util.Set;
/**
* HanaDialect
*
* @Author: cl1226
* @Date: 2025/3/14 17:47
**/
public class HanaDialect extends AbstractDialect {
private static final int MAX_DECIMAL_PRECISION = 38;
private static final int MIN_DECIMAL_PRECISION = 1;
@Override
public Set<LogicalTypeRoot> supportedTypes() {
return EnumSet.of(
LogicalTypeRoot.CHAR,
LogicalTypeRoot.VARCHAR,
LogicalTypeRoot.BOOLEAN,
LogicalTypeRoot.VARBINARY,
LogicalTypeRoot.DECIMAL,
LogicalTypeRoot.TINYINT,
LogicalTypeRoot.SMALLINT,
LogicalTypeRoot.INTEGER,
LogicalTypeRoot.BIGINT,
LogicalTypeRoot.FLOAT,
LogicalTypeRoot.DOUBLE,
LogicalTypeRoot.DATE,
LogicalTypeRoot.TIME_WITHOUT_TIME_ZONE,
LogicalTypeRoot.TIMESTAMP_WITHOUT_TIME_ZONE);
}
@Override
public String dialectName() {
return "hana";
}
@Override
public JdbcRowConverter getRowConverter(RowType rowType) {
return new HanaRowConverter(rowType);
}
@Override
public Optional<String> defaultDriverName() {
return Optional.of("com.sap.db.jdbc.Driver");
}
@Override
public String getLimitClause(long limit) {
return "LIMIT " + limit;
}
@Override
public String quoteIdentifier(String identifier) {
return "\"" + identifier + "\"";
}
@Override
public Optional<String> getUpsertStatement(String tableName, String[] fieldNames, String[] uniqueKeyFields) {
return Optional.empty();
}
@Override
public Optional<Range> decimalPrecisionRange() {
return Optional.of(Range.of(MIN_DECIMAL_PRECISION, MAX_DECIMAL_PRECISION));
}
}
HanaDialectFactory
package org.apache.flink.connector.jdbc.databases.hana.dialect;
import org.apache.flink.annotation.Internal;
import org.apache.flink.connector.jdbc.dialect.JdbcDialect;
import org.apache.flink.connector.jdbc.dialect.JdbcDialectFactory;
/** Factory for {@link HanaDialect}. */
@Internal
public class HanaDialectFactory implements JdbcDialectFactory {
@Override
public boolean acceptsURL(String url) {
return url.startsWith("jdbc:sap:");
}
@Override
public JdbcDialect create() {
return new HanaDialect();
}
}
HanaRowConverter
package org.apache.flink.connector.jdbc.databases.hana.dialect;
import org.apache.flink.annotation.Internal;
import org.apache.flink.connector.jdbc.converter.AbstractJdbcRowConverter;
import org.apache.flink.table.types.logical.RowType;
/**
* Runtime converter that responsible to convert between JDBC object and Flink internal object for
* MySQL.
*/
@Internal
public class HanaRowConverter extends AbstractJdbcRowConverter {
private static final long serialVersionUID = 1L;
@Override
public String converterName() {
return "Hana";
}
public HanaRowConverter(RowType rowType) {
super(rowType);
}
}
2. META-INF\services中增加
org.apache.flink.connector.jdbc.databases.hana.dialect.HanaDialectFactory
3. 使用方式
drop table if exists sapserver_ods_xxx_f;
create table sapserver_ods_xxx_f(
MANDT string ,
LAENG decimal(13,3) ,
BREIT decimal(13,3) ,
MEABM string ,
VOLUM decimal(13,3) ,
VOLEH string ,
BRGEW decimal(13,3) ,
TOP_LOAD_FULL decimal(15,3) ,
TOP_LOAD_FULL_UOM string ,
CAPAUSE decimal(15,3) ,
TY2TQ string,
DUMMY_UOM_INCL_EEW_PS string ,
`/CWM/TY2TQ` string ,
`/STTPEC/NCODE` string,
`/STTPEC/NCODE_TY` string,
PCBUT string
) with (
'connector' = 'jdbc',
${sap_hana},
'table-name' = 'xxx',
'driver' = 'com.sap.db.jdbc.Driver'
);
更多推荐
所有评论(0)