生产中遇到一些譬如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'

);

Logo

腾讯云面向开发者汇聚海量精品云计算使用和开发经验,营造开放的云计算技术生态圈。

更多推荐