项目需求为DM8DSC数据实时同步到kafka中,在同步过程中还需要满足以下几点要求:
1.消息在异常场景不丢失,保持顺序;
2.KEY值相同的数据,发送到相同的kafka topic partition;
3.需要按照特定的json格式发送消息;
4.源端带并发时(例如并发200),延迟不超过2s;
DMHS 支持将源端数据库的数据以 json 格式发送至 kafka 进行消费,当同步目标为 kafka 消息队列时,DMHS 执行端以动态库的方式通过 JNI 接口被 java 同步服务程序调用,然后 java 同步服务程序将获取的 json 消息发送给 kafka端进行消费同步,如图所示:
1)检查是否开启归档:
select arch_mode from v$database;
注:ARCH_MODE 为“Y”,表示启用归档;为“N”,表示未启用。
2)检查归档配置位置:
cat dmarch.ini
注:DM8DSC归档可以放在ASM上,也可以放在本地盘上。
1)检查是否打开逻辑日志:
select para_name,para_value from v$dm_ini where para_name='RLOG_APPEND_LOGIC';
注:1或者2,代表已打开逻辑日志,其中:
1:如果有主键列,记录 UPDATE 和 DELETE 操作时只包含主键列信息,若没有主键列则包含所有列信息;
2:不论是否有主键列,记录 UPDATE 和 DELETE 操作时都包含所有列的信息;
如果传入到kafka的消息中,针对update、delete操作需要包含全列信息,则需要设置为2,不需要可以设置为1。建议设置为2.
2)打开逻辑日志:
sp_set_para_value(1, 'RLOG_APPEND_LOGIC',2);
#动态参数,不需要重启就能生效
HS 12月月度版后,源端为DSC的场景需要将DSC的DSC_TRX_VIEW_SYNC参数设置为1才能正常启动HS服务,否则报错:
sp_set_para_value(1, 'DSC_TRX_VIEW_SYNC',1);
#动态参数,不需要重启就能生效。新版本均建议DSC环境将该参数设置为1.
1)上传软件到服务器后,更改dmhs软件所属组和用户,创建dmhs目录:
(源端DSC使用dmdba安装,目的端kafka使用root安装,以下以DSC端为例)
chown dmdba.dinstall dmhs_V4.3.44_dm8_kafka_rev229216_FTarm_kylin4_64_20260803.bin
chmod 755 dmhs_V4.3.44_dm8_kafka_rev229216_FTarm_kylin4_64_20260803.bin
2)安装DMHS:
[dmdba@localhost tmp]$ ./dmhs_V4.3.44_dm8_kafka_rev229216_FTarm_kylin4_64_20260803.bin -i
Extract install files..........
1.英文(English)
2.简体中文(简体中文)
请选择安装语言[2.简体中文(简体中文)]:2
/tmp/DMHSInstall/install.log
1.免费试用达梦数据实时同步
2.使用已申请的Key文件
验证许可证文件[1.免费试用达梦数据实时同步]:
1.精简版
2.完整版(web客户端)
3.自定义
安装类型[1.精简版]:2
1.实时同步软件服务器
2.远程部署工具
3.实时同步软件客户端
4.内置数据库
5.实时同步软件配置助手
6.手册
所需磁盘空间:900 MB
安装目录: [/home/dmdba/dmhs] /data/dmhs/dmhs/
该路径不为空,是否继续安装?[Y or N]Y
安装路径可能存在覆盖安装
1.统一部署
2.现在初始化
是否初始化达梦数据实时同步系统[1.统一部署]:1
正在安装
default start ... default finished.
server start ... server finished.
hs_agent start ... hs_agent finished.
webmanager start ... webmanager finished.
db start ... db finished.
hsca start ... hsca finished.
doc start ... doc finished.
doc start ... doc finished.
postinstall start ... postinstall finished.
正在创建快捷方式
安装成功
远程部署工具配置
远程部署工具名称[HsAgent]:
主机Ip(外网)[192.168.122.1](192.168.122.1,192.168.56.66):10.20.xx.xx
远程部署工具管理端口[5456](1000-65535):
内置数据库轮询间隔[3](1-60):
内置数据库IP[192.168.122.1]:
内置数据库端口[15236]:
内置数据库用户名[SYSDBA]:
内置数据库密码[SYSDBA]:
远程控制服务
1.自动
2.手动
启动方式:[2.手动]
正在创建远程控制服务
web服务
1.自动
2.手动
启动方式:[2.手动]
正在创建web服务
达梦数据实时同步V4.0安装完成
更多安装信息,请查看安装日志文件:
/home/dmdba/dmhs/log/install.log
[dmdba@localhost tmp]$
更多安装信息,请查看安装日志文件:
/opt/dmhs/log/install.log
当HS连接DM数据库时,需要使用libdmoci.so文件,所以需要申请dmoci包(注:dmoci版本需要与DM数据库版本一致)。然后执行以下命令将libdmoci.so放到源端DSC数据库bin目录下与hs bin2目录下:
unzip **_ent_8.1.4.80_dmdci.zip
cp dmocci/libdmoci.so /home/dmdba/dmdbms/bin
cp dmocci/libdmoci.so /data/dmhs/dmhs/bin2
上传dmhs.key到bin 目录,修改所属组和用户。
chown dmdba. ..../bin/dmhs.key
检查需要部署dmhs的机器上的环境变量是否已包含dm8数据库的bin目录,如果没有,则添加。例如:假设dm8安装路径为/opt/dm8
1)执行vim ~/.bash_profile
2)在文件中加入export LD_LIBRARY_PATH=/opt/dm8/bin
3)保存修改,然后执行source ~/.bash_profile使其生效
捕获端、执行端数据库用户必须具有同步对象的操作权限及操作用户下建表的权限。本次数据源端同步用户采用SYSDBA用户。若不允许使用SYSDBA用户,则可单独创建同步用户DMHS,并授予以下权限:
#创建用户
Create user DMHS identified by “XXXXX”;
#授权
sf_set_system_para_value('ENABLE_DDL_ANY_PRIV',0,0,1);
grant "PUBLIC","RESOURCE","SOI","VTI" to "DMHS";
grant CREATE ROLE to "DMHS";
grant CREATE TABLE to "DMHS";
grant CREATE VIEW to "DMHS";
grant CREATE PROCEDURE to "DMHS";
grant CREATE SEQUENCE to "DMHS";
grant CREATE TRIGGER to "DMHS";
grant CREATE INDEX to "DMHS";
grant CREATE CONTEXT INDEX to "DMHS";
grant CREATE LINK to "DMHS";
grant CREATE REPLICATE to "DMHS";
grant CREATE PACKAGE to "DMHS";
grant CREATE SYNONYM to "DMHS";
grant CREATE PUBLIC SYNONYM to "DMHS";
grant ALTER REPLICATE to "DMHS";
grant DROP REPLICATE to "DMHS";
grant DROP ROLE to "DMHS";
grant ADMIN ANY ROLE to "DMHS";
grant ADMIN ANY DATABASE PRIVILEGE to "DMHS";
grant CREATE ANY SCHEMA to "DMHS";
grant DROP ANY SCHEMA to "DMHS";
grant CREATE ANY TABLE to "DMHS";
grant ALTER ANY TABLE to "DMHS";
grant DROP ANY TABLE to "DMHS";
grant INSERT TABLE to "DMHS";
grant INSERT ANY TABLE to "DMHS";
grant UPDATE TABLE to "DMHS";
grant UPDATE ANY TABLE to "DMHS";
grant DELETE TABLE to "DMHS";
grant DELETE ANY TABLE to "DMHS";
grant SELECT TABLE to "DMHS";
grant SELECT ANY TABLE to "DMHS";
grant REFERENCES TABLE to "DMHS";
grant REFERENCES ANY TABLE to "DMHS";
grant CREATE ANY VIEW to "DMHS";
grant ALTER ANY VIEW to "DMHS";
grant DROP ANY VIEW to "DMHS";
grant INSERT VIEW to "DMHS";
grant INSERT ANY VIEW to "DMHS";
grant UPDATE VIEW to "DMHS";
grant UPDATE ANY VIEW to "DMHS";
grant DELETE VIEW to "DMHS";
grant DELETE ANY VIEW to "DMHS";
grant SELECT VIEW to "DMHS";
grant SELECT ANY VIEW to "DMHS";
grant CREATE ANY PROCEDURE to "DMHS";
grant DROP ANY PROCEDURE to "DMHS";
grant EXECUTE PROCEDURE to "DMHS";
grant EXECUTE ANY PROCEDURE to "DMHS";
grant CREATE ANY SEQUENCE to "DMHS";
grant DROP ANY SEQUENCE to "DMHS";
grant SELECT SEQUENCE to "DMHS";
grant SELECT ANY SEQUENCE to "DMHS";
grant CREATE ANY TRIGGER to "DMHS";
grant DROP ANY TRIGGER to "DMHS";
grant CREATE ANY INDEX to "DMHS";
grant ALTER ANY INDEX to "DMHS";
grant DROP ANY INDEX to "DMHS";
grant CREATE ANY CONTEXT INDEX to "DMHS";
grant ALTER ANY CONTEXT INDEX to "DMHS";
grant DROP ANY CONTEXT INDEX to "DMHS";
grant CREATE ANY PACKAGE to "DMHS";
grant DROP ANY PACKAGE to "DMHS";
grant EXECUTE PACKAGE to "DMHS";
grant EXECUTE ANY PACKAGE to "DMHS";
grant CREATE ANY LINK to "DMHS";
grant DROP ANY LINK to "DMHS";
grant CREATE ANY SYNONYM to "DMHS";
grant DROP ANY SYNONYM to "DMHS";
grant DROP PUBLIC SYNONYM to "DMHS";
grant SELECT ANY DICTIONARY to "DMHS";
grant ADMIN REPLAY to "DMHS";
grant ADMIN BUFFER to "DMHS";
grant ALTER ANY TRIGGER to "DMHS";
grant CREATE MATERIALIZED VIEW to "DMHS";
grant CREATE ANY MATERIALIZED VIEW to "DMHS";
grant DROP ANY MATERIALIZED VIEW to "DMHS";
grant ALTER ANY MATERIALIZED VIEW to "DMHS";
grant SELECT MATERIALIZED VIEW to "DMHS";
grant SELECT ANY MATERIALIZED VIEW to "DMHS";
grant CREATE ANY DOMAIN to "DMHS";
grant DROP ANY DOMAIN to "DMHS";
grant CREATE DOMAIN to "DMHS";
grant USAGE ANY DOMAIN to "DMHS";
grant USAGE DOMAIN to "DMHS";
grant CREATE ANY CONTEXT to "DMHS";
grant DROP ANY CONTEXT to "DMHS";
grant COMMENT ANY TABLE to "DMHS";
grant DUMP ANY TABLE to "DMHS";
grant DUMP TABLE to "DMHS";
grant CREATE ANY DIRECTORY to "DMHS";
grant DROP ANY DIRECTORY to "DMHS";
grant ALTER ANY SEQUENCE to "DMHS";
需要在源端数据库以SYSDBA用户,在SYSDBA模式下创建 DDL 触发器及 DDL 记录表,详细参照 ddl_sql_dm8.sql脚本。(DM必须使用SYSDBA用户创建)
<?xml version="1.0" encoding="GB2312" standalone="no"?>
<dmhs>
<base>
<lang>ch</lang>
<mgr_port>5345</mgr_port>
<name>hs_cpt_dm</name>
<ckpt_interval>60</ckpt_interval>
<siteid>1</siteid>
<version>2.0</version>
</base>
<cpt>
<name>cpt0</name>
<db_type>DSC</db_type>
<db_server>**.**.0.122</db_server>
<db_user>SYSDBA</db_user>
<db_pwd>SYSDBA</db_pwd>
<char_code>PG_UTF8</char_code>
<is_ucvt>1</is_ucvt>
#源端dsc字符集为GB18030,但由于测试过程中中文一直乱码,所以将char_code调整为utf-8,并且增加is_ucvt参数对解析的日志进行转码,解决乱码问题
<db_port>5237</db_port>
<ddl_mask>0</ddl_mask>
<parse_thr>6</parse_thr> #日志分析线程数,多线程解析日志
<arch>
<clear_flag>0</clear_flag>
<clear_interval>600</clear_interval>
</arch>
<dm8_rac>
<nodes>2</nodes>
<rac_type>1</rac_type>
<db_server>**.**.0.122</db_server>
<db_port>9349</db_port> #与dmdcr_cfg.ini中的ASM模块的DCR_EP_PORT保持一致
<db_user>default</db_user>
<db_pwd>default</db_pwd>
<dir_replace>
<item>0#+DMDARCH/DSC0/arch</item>
<item>1#+DMDARCH/DSC1/arch</item>
</dir_replace>
<epoch>3</epoch>
</dm8_rac>
<send>
<ip>**.**.0.130</ip>
<mgr_port>5347</mgr_port>
<net_pack_size>256</net_pack_size>
<data_port>5348</data_port>
<trigger>0</trigger>
<constraint>0</constraint>
<identity>0</identity>
<net_turns>0</net_turns>
<crc_check>0</crc_check>
<filter> #投递消息过滤
<enable> <item>模式名.表名</item>
</enable>
</filter>
<map/>
</send>
</cpt>
</dmhs>
注:源端DSC为大小写不敏感库,
<?xml version="1.0" encoding="GB2312" standalone="no"?>
<dmhs>
<base>
<lang>ch</lang>
<mgr_port>5347</mgr_port>
<chk_interval>3</chk_interval>
<ckpt_interval>60</ckpt_interval>
<siteid>2</siteid>
<version>2.0</version>
</base>
<exec>
<recv>
<data_port>5348</data_port>
</recv>
<name>yondif_datasync</name>
<enable>1</enable>
<char_code>PG_GB18030</char_code>
<level>0</level>
<exec_thr>8</exec_thr>
<exec_policy>2</exec_policy>
<toggle_case>0</toggle_case>
<ddl_mode>1</ddl_mode>
<commit_policy>1</commit_policy>
<enable_merge>0</enable_merge>
<is_kafka>1</is_kafka>
<json_format>file</json_format> #使用自定义配置文件json_format.ini中的格式
<max_packet_size>16</max_packet_size>
</exec>
</dmhs>
dmhs_kafka.properties配置文件可以设置 kafka集群的地址、kafka 的 topic 名称、是否进行 json 格式校验、kafka 消息确认 ack 参数、kafka生产者相关优化参数等。该文件与目的端kafka的dmhs.hs在同一部目录下:
[root]# cat /data/dmhs/dmhs/bin2/dmhs_kafka.properties
# DMHS config file path
dmhs.conf.path=/data/dmhs/dmhs/bin2/dmhs.hs
# kafka broker list,such as ip1:port1,ip2:port2,...
bootstrap.servers=**.**.0.130:34999
# kafka topic name
kafka.topic.name=yondif_datasync
#控制发送的消息到指定的 partition 时,使用此参数指定发送 json消息中表名或者某个字段的解析格式,每层之间以英文:隔开。如果配置了此参数,kaka 生产者参数 partitioner.class 需注释,使用 kafka 生产者默认分区策略。
dmhs.sendKey.parse.format=key
# kafka partitioner config
#partitioner.class=com.dameng.dmhs.dmga.service.impl.OnePartitioner
# whether to enable JSON format check
json.format.check=1
# How many messages print cost time,生产建议设置为1000,提高性能
print.message.num=1000
# How many messages batch to get,生产建议设置为100,提高性能
dmhs.min.batch.size=100
# kafka serializer class
key.serializer=org.apache.kafka.common.serialization.StringSerializer
value.serializer=org.apache.kafka.common.serialization.StringSerializer
# kafka request acks config
acks=-1
max.request.size=5024000
retries=10000000
compression.type=none
max.in.flight.requests.per.connection=1
send.buffer.bytes=1048576
metadata.max.age.ms=300000
#batch.size=1048576
#linger.ms=3
#buffer.memory=134217728
#enable.idempotence=true
注:由于本项目要求KEY值相同的数据,发送到相同的kafka topic partition,所以在prperties文件中需要配置dmhs.sendKey.parse.format,并且将partitioner.class注释掉。本项目为kafka提供的消息格式设置为:
JSON_FORMAT_INS = {
"table":"#SCHEMA.#TABLE",
"op_type":"#OP_TYPE",
"key":"#ROWID",
"op_ts":"#OP_TIME",
"post_ts":"#TIME",
"pos":"#LSN",#CONCAT
"primary_keys":[#PRIMARY_KEY],#CUT
"data":{#NEW_VALUES}
},其中key属性封装的值为源端数据的rowid,为了保证每个top partition中数量均衡。所以需要设置dmhs.sendKey.parse.format=key
当kafka端dmhs.hs中json_format设置为file时,可以使用json_format.ini 配置文件指定输出的json格式,该文件默认在dmhs.hs同目录:
[root]# cat /data/dmhs/dmhs/bin2/json_format.ini
#DMHS kafka json format Configuration file
#this is comments
#common format control parameters
BATCH_COMMIT = 1
OP_TIME_FORMAT = (yyyy-mm-dd hh:mi:ss.fff)
CUR_TIME_FORMAT =(yyyy-mm-dd hh:mi:ss.fff)
COMMIT_TIME_FORMAT = (yyyy-mm-dd hh:mi:ss.fff)
NEED_CRLF = 0
OPCMD_LEN = 7
SET_NULL = 1
SET_QUOTA = 0
#CHAR_REPLACE =(",\"),(\,\\),(0x0D,\n),(0x0A,\r),(0x0c,\f)
CHAR_REPLACE =(",\"),(\,\\)
ADD_TABLE_TOPIC = 0
NEW_VALUES = ALL
SET_ROWID_COL = 0
LOB_PIECE = 0
CLOB_FORMAT = CHAR
#BLOB_FORMAT = BASE64
JSON_FORMAT_INS = {"table":"#SCHEMA.#TABLE","op_type":"#OP_TYPE","key":"#ROWID","op_ts":"#OP_TIME","post_ts":"#TIME","pos":"#LSN",#CONCAT"primary_keys":[#PRIMARY_KEY],#CUT"data":{#NEW_VALUES}}
JSON_FORMAT_UPD = {"table":"#SCHEMA.#TABLE","op_type":"#OP_TYPE","key":"#ROWID","op_ts":"#OP_TIME","post_ts":"#TIME","pos":"#LSN",#CONCAT"primary_keys":[#PRIMARY_KEY],#CUT"data":{#NEW_VALUES}}
JSON_FORMAT_DEL = {"table":"#SCHEMA.#TABLE","op_type":"#OP_TYPE","key":"#ROWID","op_ts":"#OP_TIME","post_ts":"#TIME","pos":"#LSN",#CONCAT"primary_keys":[#PRIMARY_KEY],#CUT"data":{#OLD_VALUES}}
注:新kafka-json格式将INSERT|UPDATE|DELETE|DDL拆分成多个配置项,根据操作类型,按照对应的输出格式进行输出。如果JSON不需要换行符,可以按照以上格式进行书写,将json写到一行中。如果需要换行,可以按照以下格式书写:
JSON_FORMAT_INS = {
"table":"#SCHEMA.#TABLE",
"op_type":"#OP_TYPE",
"key":"#ROWID",
"op_ts":"#OP_TIME",
"post_ts":"#TIME",
"pos":"#LSN",#CONCAT
"primary_keys":[#PRIMARY_KEY],#CUT
"data":{#NEW_VALUES}
}
当需要以表名作为 topic 发送时,可在tableTopicMap.properties中配置表名与topic的映射规则:
[root]$ cat /data/dmhs/dmhs/bin2/tableTopicMap.properties
# TABLE1=TOPIC1
DM_TEST=DM_TEST
[dmdba bin2]$ cp TemplateDmhsService DmhsServiceHSSERVER
[dmdba bin2]$ vi DmhsServiceHSSERVER
#REPLACE DMHS_HOME path
DMHS_HOME=/data/dmhs/dmhs
#REPLACE program dir
PROG_DIR=/data/dmhs/dmhs/bin2
#REPLACE program config path
CONF_PATH=/data/dmhs/dmhs/bin2/dmhs.hs
#REPLACE need library path, LD_LIBRARY_PATH/LIBPATH
NEED_LIB_PATH=
HS_NLS_LANG=zh_CN.UTF-8
注:为了解决乱码问题,HS_NLS_LANG需要配置与源端dmhs.hs中char_code相同的字符集
[root bin2]# vi start_dmhs_kafka.sh
#!/bin/bash
export LANG=zh_CN.UTF-8 #与源端字符集保持一致,避免乱码
/data/iuap/middleware/openjdk8/bin/java -Djava.ext.dirs="/data/iuap/middleware/kafka-34999/libs:." com.dameng.dmhs.dmga.service.impl.ExecDMHSKafkaService /data/dmhs/dmhs/bin2/dmhs_kafka.properties
[root]$ cd /data/dmhs/dmhs/bin2
[root]$ nohup ./start_dmhs_kafka.sh &
[root]$./dmhs_console
DMHS>start exec
#观察控制台及日志,是否反馈执行成功。
#日志:tail -200f log/dmhs_202606.log
[dmdba@localhost ~]$ cd /data/dmhs/dmhs/bin2
[dmdba@localhost bin2]$ ./DmhsServiceHSSERVER start
#查看日志是否启动成功:tail -200f log/dmhs_202606.log
[dmdba@localhost bin2]$ ./dmhs_console
#清理目的端检查点,首次启动HS时执行,一旦开始装载及同步后就不能再执行了:
DMHS>clear exec lsn
#针对需要同步的表,装载字典(针对特定表装载字典,非同步表不装载,解析归档时可以只解析需要同步的表,降低解析归档压力):
DMHS> copy 0"sch.name='xxxx' AND TAB.NAME in('xxxx')" CLEAR|DICT
CPT模块在全量数据装载完成,需要增量同步时开启。全量装载时需要保证从装载开始到装载结束,源端归档是完整的。
[dmdba@localhost bin2]$ ./dmhs_console
DMHS>start cpt
#观察控制台及日志,是否反馈执行成功。
#日志:tail -200f log/dmhs_202606.log
以xx.T1表为例:
[dmdba@localhost bin2]$ ./dmhs_console
#装载前先停止CPT模块
DMHS>stop cpt
#单表装载
DMHS>copy 0 "sch.name='xx' and tab.name='T1'" insert|extent|group|16
#按模式装载
DMHS>copy 0 "sch.name='xx'" insert|extent|group|16|thread|8
#观察控制台及日志,是否装载完成
#日志:tail -200f log/dmhs_202606.log
#装载完成后,开启CPT模块,打开增量同步
DMHS>start cpt
#kafka 端执行以下操作可以观察到数据接收情况:
[root]$ cd /data/iuap/middleware/kafka-34999/bin
#查看消费者情况,打开后可以观察到json消息:
[root]$ ./kafka-console-consumer.sh --bootstrap-server localhost:34999 --topic yondif_datasync
#以下命令可以观察到partition情况:
[root]$ ./kafka-console-consumer.sh --bootstrap-server localhost:34999 --topic yondif_datasync --property print.key=true --property print.value=true --property print.timestamp=true --property print.partition=true
#源端登录数据库执行操作并观察是否同步到kafka:
update xx.T1 set ELES_ID='3' where BALANCE_ID='811859214686224385';commit;
发送到kafka消息如下:
文章
阅读量
获赞
