注册
达梦数据复制软件DMDRS同步配置文件模板:以Kafka为目标数据库
专栏/数+融合专栏/ 文章详情 /

达梦数据复制软件DMDRS同步配置文件模板:以Kafka为目标数据库

达梦数+数据融合软件(DMDF)产品团队 2026/07/30 159 0 0
摘要 介绍目标数据库为kafka时,目标DMDRS的配置文件模板。

DMDRS使用XML文件作为配置文件,通过XML配置文件配置产品服务功能。常用产品服务的配置文件模板存放于DMDRS安装目录下,配置文件模板名称和默认配置文件名称如下表所示。

说明

启动DMDRS服务时,如不指定配置文件,缺省使用默认配置文件。

DMDRS服务 配置文件模板 默认配置文件
源DMDRS cpt.xml、cpt_dsc.xml drs.xml
目标DMDRS exec_kafka.xml drs.xml
DMDSS dss.xml dss.xml

DMDRS的XML配置文件的格式如下表所示。

组成部分 语法 说明
声明语句 <?xml version="1.0" encoding="GB18030"?> 包括XML版本号和字符集声明。
根元素 <drs></drs> 根元素中包括所有的子元素。
子元素 Manager模块:<base></base>
CPT模块:<cpt></cpt>
DSS模块:<dss></dss>
EXEC模块:<exec></exec>
各模块对应的子元素中还包括具体的元素,具体请参见《DMDRS参考手册》

1 配置源DMDRS

源DMDRS服务包括Manager模块和CPT模块。在源DMDRS服务中,Manager模块提供基本的服务管理功能,CPT模块提供数据捕获功能,模块的详细介绍请参见《DMDRS产品介绍》。各源数据库类型的配置文件模板参见源数据库对应的DMDRS搭建手册。

2 配置目标DMDRS

目标DMDRS服务包括Manager模块和EXEC模块。在目标DMDRS服务中,Manager模块提供基本的服务管理功能,EXEC模块提供数据执行功能,模块的详细介绍请参见《DMDRS产品介绍》

配置文件模板

<?xml version="1.0" encoding="GB18030"?> <drs> <base> <mgr_port>Manager模块的端口号</mgr_port> <siteid>站点号</siteid> </base> <exec> <name>EXEC模块的名称</name> <login> <dbtype>目标数据库的类型</dbtype> <producer_properties>Kafka生产者属性文件路径</producer_properties> </login> <group> <item> <id>分组编号</id> <topic_name>分组的目标Kafka Topic名称</topic_name> <desc> <table>分组的表匹配信息</table> </desc> <json_format_path>JSON模板文件路径</json_format_path> </item> </group> </exec> </drs>

参数说明

配置文件模板中的配置参数为必配参数,包括Manager模块的端口号、站点号、EXEC模块的名称和目标数据库的登录信息等。其余选配参数由用户根据实际功能或性能需求选择性地配置,参数配置请参见《DMDRS参考手册》

  • producer_properties参数配置所需要的Kafka生产者属性文件模板请参见Kafka生产者属性配置

  • json_format_path参数配置所需要的JSON模板请参见JSON模板配置

  • 如EXEC模块采用DSS模块消费模式,需配置dss参数。

  • 如需提高数据执行入库的效率,可配置group参数。

    如上示例中配置参数为必配参数,其余选配参数请参见《DMDRS参考手册》

2.1 Kafka生产者属性配置

Kafka生产者属性是目标DMDRS服务使用librdkafka驱动与Kafka服务建立连接的基本属性,如连接的目标IP和端口,连接使用的账号和密码,以及一些性能相关参数等。需要在目标DMDRS的XML配置文件中配置producer_properties参数为Kafka生产者属性配置文件路径。

生产者属性文件模板

  • 基本配置模板

    bootstrap.servers=Kafka节点1的IP:节点1的端口号<,Kafka节点2的IP:节点2的端口号...>
    
  • Kafka配置了SASL认证下的配置模板

    bootstrap.servers=Kafka节点1的IP:节点1的端口号<,Kafka节点2的IP:节点2的端口号...>
    security.protocol=安全协议名 #如SASL_PLAINTEXT、SASL_SSL
    sasl.mechanism=安全策略名    #如PLAIN, SCRAM-SHA-256, SCRAM-SHA-512
    sasl.username=SASL用户名
    sasl.password=SASL密码
    
  • Kafka配置了Kerberos认证的配置模板

    bootstrap.servers=Kafka节点1的IP:节点1的端口号<,Kafka节点2的IP:节点2的端口号...>
    security.protocol=SASL_PLAINTEXT
    sasl.mechanism=GSSAPI
    sasl.kerberos.principal=Kerberos主体                #如dameng/HADOOP.COM@HADOOP.COM
    sasl.kerberos.keytab=keytab文件路径                                                     
    sasl.kerberos.service.name=Kafka服务的kerberos主体   #如kafka
    

参数说明

生产者属性文件模板中的配置参数为必配参数,包括Kafka服务的IP、端口号和登录信息。其余选配参数由用户根据实际功能或性能需求选择性地配置,生产者属性说明请参见生产者属性列表

常用性能参数

librdkafka驱动包含许多影响数据同步性能的生产者属性。常用的性能相关的生产者属性示例如下:

说明

以下DMDRS服务的配置文件仅为示例,在实际应用场景中,请根据实际搭建环境修改配置文件中的参数配置,生产者属性说明请参见生产者属性列表

acks=-1
message.max.bytes=5024000
retries=65535
compression.codec=none
max.in.flight.requests.per.connection=1
socket.send.buffer.bytes=1048576
metadata.max.age.ms=300000
batch.size=524288
linger.ms=20

2.2 JSON模板配置

目标DMDRS根据JSON模板文件控制日志操作转换为JSON或其它格式文本串的打印内容与格式。JSON模板文件内容主要包括控制参数、变量参数、表达式变量的声明和使用,具体请参见《DMDRS参考手册》json_format_path参数说明。

JSON模板配置示例如下:

  • 基础JSON模板

    #DMDRS Kafka JSON format configuration file
    #This is a comment
    
    #Common format control parameters
    OP_TIME_FORMAT                   = (yyyy-mm-dd hh:mi:ss)
    CUR_TIME_FORMAT                  = (yyyy-mm-dd hh:mi:ss)
    NULL_FORMAT                      = "null"
    SET_QUOTA                        = 0
    SET_ROWID_COL                    = 0
    NEED_CRLF                        = 1
    CHAR_REPLACE                     = (\,\\),(",\"),(0x0D,\r),(0x0A,\n),(0x0c,\f),(/,/\),(0x08,\b),(0x09,\t)
    NEW_VALUES                       = ALL
    OLD_LOB_FLAG                     = "empty_lob()"
    JSON_FORMAT_INS = {
        "table":"#SCHEMA.#TABLE",
        "op_type":"#OP_TYPE",
        "op_ts":"#OP_TIME",
        "current_ts":"#TIME",
        "pos":"#POS",
        "primary_keys":[#PRIMARY_KEY],
        "after":{#NEW_VALUES}
    }
    JSON_FORMAT_UPD = {
        "table":"#SCHEMA.#TABLE",
        "op_type":"#OP_TYPE",
        "op_ts":"#OP_TIME",
        "current_ts":"#TIME",
        "pos":"#POS",
        "primary_keys":[#PRIMARY_KEY],
        "before":{#OLD_VALUES},
        "after":{#NEW_VALUES}
    }
    JSON_FORMAT_DEL = {
        "table":"#SCHEMA.#TABLE",
        "op_type":"#OP_TYPE",
        "op_ts":"#OP_TIME",
        "current_ts":"#TIME",
        "pos":"#POS",
        "primary_keys":[#PRIMARY_KEY],
        "before":{#OLD_VALUES}
    }
    JSON_FORMAT_DDL = {
        "table":"#SCHEMA.#TABLE",
        "op_type":"#OP_TYPE",
        "op_ts":"#OP_TIME",
        "current_ts":"#TIME",
        "pos":"#POS",
        "ddl_text":"#DDL_SQL"
    }
    
  • 兼容华为DRS的JSON模板

    #DMDRS kafka json format Configuration file
    #This is a comment
    
    #common format control parameters
                    OP_TIME_FORMAT                                  = (yyyy-mm-dd hh:mi:ss.ff)
                    CUR_TIME_FORMAT                                 = (yyyy-mm-ddThh:mi:ss.ff)
                    NEED_CRLF                                       = 1
                    SET_QUOTA                                       = 1
                    CHAR_REPLACE                                    = (\,\\),(",\"),(0x0D,\r),(0x0A,\n),(0x0c,\f),(/,/\),(0x08,\b),(0x09,\t)
                    NEW_VALUES                                      = ALL
                    OLD_VALUES                                      = ALL
                    SET_ROWID_COL                                   = 0
                    SET_TIME_EPOCH                                  = 2
                    NULL_FORMAT                                     = ""
                    OLD_LOB_FLAG                                    = "empty_lob()"
                    UPDATE_FILL_FLAG                                = 0
                    PARTITION_NO                                    = -1
                    SET_DEFAULT_KEY                                 = 2 
                    MAX_LOB_SIZE                                    = 0
                    NAME_TYPE_CASE                                  = 1
                    BLOB_FORMAT                                     = NUM_ARRAY
        
    JSON_FORMAT= {
        "schema":"#SCHEMA",
        "dbtype":"#DBTYPE",
        "type":"#OP_TYPE",
        "op_type":"#OP_TYPE",
        "pkNames":[#PRIMARY_KEY],
        "sql":"",
        "table":"#TABLE",
        "is_ddl":false,
        "es":#OP_TIME,
        "ts":#TIME,
        "id":#POS,
        "columnType":{
    #NAME_TYPE 
        },
        "sqlType": {
    #SQLTYPE
        },
        "data":#CONCAT[
            {
    #NEW_VALUES
            }
        ]#CUT_NULL,
        "old":#CONCAT[
            {
    #OLD_VALUES
            }
        ]#CUT_NULL,
        "pk":#CONCAT{
    #PK_VALUES
        }#CUT_NULL
    }
    
    JSON_FORMAT_DDL ={
        "schema":"#SCHEMA",
        "dbtype":"#DBTYPE",
        "type":"#OP_TYPE",
        "op_type":"#OP_TYPE",
        "pkNames":[#PRIMARY_KEY],
        "sql":"#DDL_SQL",
        "table":"#TABLE",
        "is_ddl":true,
        "es":#OP_TIME,
        "ts":#TIME,
        "id":#POS,
        "columnType":{
    #NAME_TYPE 
        },
        "sqlType": {
    #SQLTYPE
        },
        "data":#CONCAT[
            {
    #NEW_VALUES
            }
        ]#CUT_NULL,
        "old":#CONCAT[
            {
    #OLD_VALUES
            }
        ]#CUT_NULL,
        "pk":#CONCAT{
    #PK_VALUES
        }#CUT_NULL
    }
    
  • 同Debezium的JSON模板

    #DMDRS kafka json format Configuration file
    #This is a comment
    
    #common format control parameters
                    OP_TIME_FORMAT                                 = (yyyy-mm-dd hh:mi:ss.ff)
                    CUR_TIME_FORMAT                                = (yyyy-mm-ddThh:mi:ss.ff)
                    NEED_CRLF                                      = 0
                    SET_QUOTA                                      = 0
                    CHAR_REPLACE                                   = (\,\\),(",\"),(0x0D,\r),(0x0A,\n),(0x0c,\f),(/,/\),(0x08,\b),(0x09,\t)
                    NEW_VALUES                                     = ALL
                    OLD_VALUES                                     = ALL
                    SET_ROWID_COL                                  = 1
                    SET_TIME_EPOCH                                 = 0
                    NULL_FORMAT                                    = ""
                    OLD_LOB_FLAG                                   = "empty_lob()"
                    UPDATE_FILL_FLAG                               = 0
                    PARTITION_NO                                   = -1
                    SET_DEFAULT_KEY                                = 2  
    
    #声明列,至少包含一个名为VALUE的值
    column version(JSONTYPE, NULLABLE, VALUE) = {$DRS_VERSION.JSON_TYPE, $DRS_VERSION.NULLABLE, $DRS_VERSION.VALUE}
    column connector(JSONTYPE, NULLABLE, VALUE) = {$DBTYPE.JSON_TYPE, $DBTYPE.NULLABLE, $DBTYPE.VALUE}
    column name(JSONTYPE, NULLABLE, VALUE) = {$SITENAME.JSON_TYPE, $SITENAME.NULLABLE, $SITENAME.VALUE}
    column ts_ms(JSONTYPE, NULLABLE, VALUE) = {$OP_TIME.JSON_TYPE, $OP_TIME.NULLABLE, $OP_TIME.VALUE}
    column snapshot(JSONTYPE, NULLABLE, VALUE) = {'boolean', 'true', 'false'}
    column db(JSONTYPE, NULLABLE, VALUE) = {$SCHEMA.JSON_TYPE, $SCHEMA.NULLABLE, $SCHEMA.VALUE}
    column table(JSONTYPE, NULLABLE, VALUE) = {$TABLE.JSON_TYPE, $TABLE.NULLABLE, $TABLE.VALUE}
    column QUERY(JSONTYPE, NULLABLE, VALUE) = {$SQL.JSONTYPE, $SQL.NULLABLE, $SQL.VALUE}
    
    #声明一段列表
    columns source = {$version, $connector, $name, $ts_ms, $snapshot, $db, $table, $LSN, $POS, $QUERY}
    
    #JSON模板
    JSON_FORMAT = {
        "schema": {
        "type": "struct",
        "fields": [
            {
            "type": "struct",
            "fields": #CONCAT[
                $NR_COL({"type": "#JSONTYPE", "optional": #NULLABLE, "field": "#NAME"})
            ]#CUT_NULL,
            "optional": true,
            "name": "#SITENAME.#SCHEMA.#TABLE.Value",  
            "field": "before"  
            },
            {
            "type": "struct",
            "fields": #CONCAT[
                $NR_COL({"type": "#JSONTYPE", "optional": #NULLABLE, "field": "#NAME"})
            ]#CUT_NULL,
            "optional": true,
            "name": "#SITENAME.#SCHEMA.#TABLE.Value",
            "field": "after"  
            },
            {
            "type": "struct",
            "fields": [ 
                $source({"type": "#JSONTYPE", "optional": #NULLABLE, "field": "#NAME"})
            ],
            "optional": false,
            "name": "io.debezium.connector.#DBTYPE.Source",  
            "field": "source"
            },
            {
            "type": "string",
            "optional": false,
            "field": "op"
            },
            {
            "type": "int64",
            "optional": true,
            "field": "ts_ms" 
            }
        ],
        "optional": false,
        "name": "#SITENAME.#SCHEMA.#TABLE.Envelope"   
        },
        "payload": {
        "op": "#OP_TYPE",
        "ts_ms": "#TIME",
        "before": #CONCAT{$OLD_COL("#NAME":#VALUE)}#CUT_NULL,
        "after": #CONCAT{$NEW_COL("#NAME":#VALUE)}#CUT_NULL,
        "sources": {
            $source({"#NAME":#VALUE})
        }
        }
    }
    
  • 兼容Oracle GoldenGate默认分隔文本格式的JSON模板

    #DMDRS kafka json format Configuration file
    #This is a comment
    
    #common format control parameters
                    SET_QUOTA                                      = 2
                    UPDATE_FILL_FLAG                               = 0
                    PARTITION_NO                                   = -1
                    SET_DEFAULT_KEY                                = 2
                    DML_TYPE_FORMAT                                = (I,U,D)
                    NON_JSON                                       = 1
    #JSON模板
    JSON_FORMAT_INS = {#OP_TYPE|#SCHEMA.#TABLE|$NEW_COL{#VALUE}(|)} 
    JSON_FORMAT_UPD = {#OP_TYPE|#SCHEMA.#TABLE|$NEW_COL{#VALUE}(|)}
    JSON_FORMAT_DEL = {#OP_TYPE|#SCHEMA.#TABLE|$OLD_COL{#VALUE}(|)}
    JSON_FORMAT_DDL = {#OP_TYPE|#SCHEMA.#TABLE}  #OGG默认不同步DDL到Kafka
    

2.3 Kafka分区配置

当需要控制发送的消息到指定的分区时,需要使用数据清洗转换(DMCVT)功能,在数据清洗转换脚本中为匹配条件的日志操作的消息配置分区key。需要在JSON_FORMAT文件中配置PARTITION_NO参数为-1。

数据清洗转换脚本的示例如下。

-- DDL配置 OBJECT "*"."*" BEGIN op.#key := CHARTOBIN(op.#OBJ); -- 示例:设置DDL消息的分区key值为对象名字符串 op.exec(); END; -- DML配置 TABLE "*"."*" BEGIN op.#key := CHARTOBIN(op.#TAB); -- 示例:设置DML消息的分区key值为表名字符串 op.exec(); END;

DMCVT功能的具体使用方法请参见《DMDRS DRS语言使用手册》

说明

可在producer.properties文件中设置partitioner参数的值为murmur2_random,以使用与Java版驱动下的默认分区器相同的hash算法计算目标分区号。

2.4 消息到主题映射

目标DMDRS支持灵活配置JSON消息到主题的映射。

映射功能 涉及的配置参数
根据模式名和表名映射到不同topic 4.3.7 group
4.3.7.1.15 topic_name
根据源端站点映射到不同topic 4.3.7.1.14 siteid
将指定模式名和表名下的消息分发到一个或多个topic 4.3.7.1.10 map
将DML和DDL消息分发到不同topic 4.3.7.1.10 map

多个映射功能之间支持叠加配置。涉及的配置参数具体说明请参见《DMDRS参考手册》

评论
后发表回复

作者

文章

阅读量

获赞

扫一扫
联系客服