
数据集成ETL大数据批处理流处理变更数据捕获【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/GitHub_Trending/se/seatunnel点击查看免费下载本教程基于 SeaTunnel 的 Zeta 引擎演示如何将 MySQL 中的客户资料表crm.customer_profile持续实时同步到 Elasticsearch并在同步过程中完成业务过滤只保留管理范围内数据、字段清洗去除手机号中的连字符、去除姓名空格、状态值映射数字状态转可读状态名以及来源元数据补充库名、表名、变更类型、同步来源标记。读完本文你将掌握一条完整的「MySQL-CDC Source → Transform 链 → Elasticsearch Sink」流水线写法理解每个配置项的底层含义并能够直接复用到自己的实时数仓或搜索索引构建场景中。流水线概览本教程要实现的完整数据链路如下MySQL-CDCSource读取crm.customer_profile表的初始快照与后续 binlog 增量变更MetadataTransform把库名Database、表名Table、行变更类型RowKind从行内元数据暴露为普通字段ReplaceTransform去掉手机号字段中的-字符SqlTransform负责行级过滤id 1000、状态映射status→ACTIVE/FROZEN/OTHER以及补充常量字段sync_sourceElasticsearchSink按主键id将数据 upsert 写入索引recipe_customer_profile。这条链路在 SeaTunnel 中完全通过一份 HOCON 配置文件声明无需编写任何 Java/Python 代码适合作为「CDC 入搜索索引」类任务的模板。前置条件在开始之前需要依次完成以下准备工作确保本地执行链路与环境就绪。1. 确认本地运行链路正常先完成 跑第一个任务确认 SeaTunnel 的本地执行链路可以正常跑通再继续本教程。2. 安装并启用所需连接器按照 部署 下载连接器插件 安装本教程需要的连接器并在config/plugin_config中保留下面两项--seatunnel-connectors-- connector-cdc-mysql connector-elasticsearch --end--然后执行安装脚本并确认两个连接器的 JAR 已经落盘cd ${SEATUNNEL_HOME} sh bin/install-plugin.sh ls connectors | grep -E connector-(cdc-mysql|elasticsearch)3. 为 Zeta 引擎准备 MySQL JDBC 驱动使用 Zeta 引擎时MySQL-CDC 连接器需要 MySQL Connector/J 驱动。把它放入${SEATUNNEL_HOME}/lib并确认 JAR 已经可见ls ${SEATUNNEL_HOME}/lib | grep mysql-connector4. 创建 CDC 专用数据库账号CDC 任务需要一个有权限读取 binlog 的 MySQL 账号。出于安全考虑不要授予修改表结构的权限例如ALTER、CREATE、DROP只授予读取与复制所需的最小权限CREATE USER IF NOT EXISTS st_user_source% IDENTIFIED BY mysqlpw; GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT, LOCK TABLES ON *.* TO st_user_source%; FLUSH PRIVILEGES;5. 检查 MySQL binlog 配置CDC 依赖 binlog 才能读取行级变更请确认以下三个变量满足要求SHOW VARIABLES WHERE variable_name IN (log_bin, binlog_format, binlog_row_image);期望值是log_bin ON、binlog_format ROW、binlog_row_image FULL。其中binlog_format ROW才能记录每一行的实际变更内容binlog_row_image FULL保证 binlog 中携带完整的前镜像/后镜像数据是 CDC 正确还原变更的基础。6. 确认 Elasticsearch 可通过 HTTPS 认证访问本示例基于 Elasticsearch 8.9.0默认启用 HTTPS 与安全认证。先用 SeaTunnel 进程可以访问的地址和账号检查连通性curl --cacert /path/to/http_ca.crt \ -u elastic:your_password https://elasticsearch.example.com:9200把可信 CA 证书复制到每个 SeaTunnel 节点并在任务配置的tls_truststore_path中填写对应的本地路径。从 ElasticsearchBaseOptions.java 的源码可以看到Elasticsearch 连接器支持tls_verify_certificate默认true、tls_verify_hostname默认true、tls_keystore_path、tls_truststore_path等一整套 TLS 配置还支持auth_typebasic/api_key/api_key_encoded多种认证方式本教程使用最基础的 Basic 认证usernamepassword。准备源表数据在 MySQL 中创建crm库和customer_profile表并插入两条初始数据注意其中一条id 900故意不在后面的管理范围内用于验证过滤逻辑CREATE DATABASE IF NOT EXISTS crm; USE crm; CREATE TABLE customer_profile ( id BIGINT NOT NULL PRIMARY KEY, name VARCHAR(64) NOT NULL, phone VARCHAR(32) NOT NULL, email VARCHAR(128) NOT NULL, status INT NOT NULL, city VARCHAR(64) NOT NULL ); INSERT INTO customer_profile (id, name, phone, email, status, city) VALUES (1001, Alice Zhang , 138-0000-1111, aliceexample.com, 1, Shanghai), (900, Bob Li, 139-8888-2222, bobexample.com, 0, Beijing);任务启动且初始快照已经可以在 Elasticsearch 中查询到之后再执行下面两条增量变更用于验证流式更新与增量插入UPDATE crm.customer_profile SET name Alice Zhang , phone 138-9999-0000, status 2 WHERE id 1001; INSERT INTO crm.customer_profile (id, name, phone, email, status, city) VALUES (1003, Carol Wang, 137-1234-8888, carolexample.com, 1, Hangzhou);完整任务配置下面给出整条链路的完整 HOCON 配置。请把示例中的主机名和凭据替换为 SeaTunnel 进程可以访问的实际值。并发运行的 MySQL CDC 任务必须使用不同的server-id范围否则会在 MySQL 集群中造成 slave ID 冲突导致任务失败。env { parallelism 1 job.mode STREAMING } source { MySQL-CDC { plugin_output mysql_customer_raw url jdbc:mysql://mysql.example.com:3306/crm username st_user_source password mysqlpw server-id 5701-5704 table-names [crm.customer_profile] startup.mode initial schema-changes.enabled false } } transform { Metadata { plugin_input mysql_customer_raw plugin_output mysql_customer_with_meta metadata_fields { Database source_database Table source_table RowKind row_kind } } Replace { plugin_input mysql_customer_with_meta plugin_output mysql_customer_cleaned replace_fields [phone] pattern - replacement is_regex false } Sql { plugin_input mysql_customer_cleaned plugin_output es_customer_profile query select id, trim(name) as name, phone, email, city, case when status 1 then ACTIVE when status 2 then FROZEN else OTHER end as status_name, source_database, source_table, row_kind, mysql_cdc as sync_source from dual where id 1000 } } sink { Elasticsearch { plugin_input es_customer_profile hosts [https://elasticsearch.example.com:9200] username elastic password elasticsearch tls_verify_certificate true tls_verify_hostname true tls_truststore_path /path/to/http_ca.crt index recipe_customer_profile primary_keys [id] max_batch_size 1 schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST data_save_mode APPEND_DATA } }把调整后的配置保存为${SEATUNNEL_HOME}/config/mysql-cdc-to-elasticsearch.conf然后使用 Zeta local 模式提交cd ${SEATUNNEL_HOME} ./bin/seatunnel.sh --config ./config/mysql-cdc-to-elasticsearch.conf -m local在另一个终端中分别在执行上面的增量 SQL 前后查询索引观察文档变化curl --cacert /path/to/http_ca.crt -u elastic:your_password \ https://elasticsearch.example.com:9200/recipe_customer_profile/_search?prettysortid配置参数源码级解读为了让读者不仅“能跑通”还能理解每个参数为什么这样写这里结合仓库源码逐一说明关键配置项。MySQL-CDC Source 关键参数参数含义源码依据server-id本连接器在 MySQL 集群中注册的从库 ID支持单个数字5400或范围5400-5408。不设置时会在6500到2148492146之间随机生成官方建议显式指定JdbcSourceOptions.javatable-names要监听的表格式为库名.表名同上文件database-names/table-names定义startup.mode启动模式可选initial/earliest/latest/specific/timestamp默认initial先做全量快照再无缝衔接增量MySqlIncrementalSourceOptions.javastop.mode停止模式可选never/latest/specific默认never持续监听增量同上文件schema-changes.enabled是否把 DDL 变更事件下发给下游。本任务下游是 Elasticsearch 索引不处理表结构变更因此设为falseSourceOptions.javastartup.mode initial意味着任务启动时会先对crm.customer_profile做一次一致性快照这也是验证阶段能立刻看到1001文档的原因快照完成后自动切换到 binlog 流式读取全过程对应用无感。Elasticsearch Sink 关键参数参数含义源码依据primary_keys用于生成文档_id的主键字段CDC 变更按主键 upsert保证更新/删除能命中同一文档ElasticsearchSinkOptions.javamax_batch_size单个 bulk 请求的最大文档数默认10。本示例显式设为1让每条变更立即可见方便观察验证同上文件MAX_BATCH_SIZE默认 10schema_save_mode索引不存在时的处理策略默认CREATE_SCHEMA_WHEN_NOT_EXIST自动建索引同上文件SCHEMA_SAVE_MODEdata_save_mode数据写入策略可选DROP_DATA/APPEND_DATA/ERROR_WHEN_DATA_EXISTS默认APPEND_DATA同上文件DATA_SAVE_MODEtls_verify_certificate/tls_verify_hostnameHTTPS 证书与主机名校验开关默认均为trueElasticsearchBaseOptions.javatls_truststore_pathPEM 或 JKS 信任库路径运行 SeaTunnel 的操作系统用户必须可读同上文件TLS_TRUST_STORE_PATH需要注意schema_save_mode与data_save_mode组合时SeaTunnel 会根据primary_keys自动决定使用 index新建文档还是 update按_id更新语义这是 CDC 场景下“更新1001后旧文档被覆盖”这一行为的实现基础。Metadata Transform 暴露的元数据字段Metadata转换用于把行内元数据提取为普通字段本示例使用了三个最常用的 CDC 元数据 Key元数据 Key输出字段名说明Databasesource_database数据所属数据库名Tablesource_table数据所属表名RowKindrow_kind行的变更类型I插入、-U更新前、U更新后、-D删除根据 metadata.md 文档MySQL-CDC 还支持EventTime、Delay、SourceTimestamp、BinlogFile、BinlogPos、BinlogRow、Gtid等更细粒度的元数据其中BinlogFile等仅在 MySQL-CDC 下有效且startup.mode initial时快照行这些字段为null。Metadata转换不会改变原有数据字段只做“追加元数据列”因此可以安全地串联在 Source 之后。验证结果任务运行后在 Elasticsearch 的_search返回结果中确认以下内容快照阶段最终只写入了1 条文档10011001被保留下来因为它的不可变主键属于id 1000的管理范围900被过滤掉因为它不在这个主键范围内trim(name)把 Alice Zhang 变成了Alice Zhang两侧空格被去除Replace把138-0000-1111变成了13800001111连字符被去除初始写入的文档里包含status_nameACTIVE、source_databasecrm、source_tablecustomer_profile、sync_sourcemysql_cdc更新1001后ES 中这条文档变成status_nameFROZEN、phone13899990000、row_kindU更新事件且旧文档被按_id覆盖插入1003后ES 中新增第二条文档字段包括status_nameACTIVE、phone13712348888、row_kindI、sync_sourcemysql_cdc插入事件。为什么这条链路这样写1.Replace负责轻量清洗在这份配置里Replace只做一件事把手机号里的-去掉。它属于轻量级字符串清洗工具按字段配置pattern、replacement并可用is_regex控制是否按正则匹配本示例is_regex false直接按字面量-匹配。把这种简单、确定性的清洗放到Replace可以让Sql专注于更复杂的业务逻辑两种 transform 各司其职。2.Sql负责业务过滤与补字段SqlTransform 的query集中了全部业务规则行过滤只保留id 1000的主键范围状态映射case when status 1 then ACTIVE when status 2 then FROZEN else OTHER end把数字状态转成可读状态名常量补字段mysql_cdc as sync_source新增来源标记字段同时透传Metadata阶段产出的source_database、source_table、row_kind字段。这里特意使用不可变主键做过滤这是一个非常重要的设计决策。不要在写入 upsert sink 之前直接用is_deleted、status这类可变字段过滤 CDC 更新如果一条已经写入 ES 的记录后来不再满足过滤条件上游不会自动为旧 ES 文档补发删除事件结果就是“该删的没删”索引里残留脏数据。对于软删除场景正确做法有两种把删除标记如is_deleted原样同步到 Elasticsearch在查询时过滤使用能将状态变化明确转换为删除事件的专用链路如依赖RowKind -D的删除事件或让下游 sink 识别删除标记并主动删除文档。3. Elasticsearch 写入配置的取舍Sink 端这组参数组合起来产生了上述预期结果primary_keys [id]以id为文档_id保证 upsert 语义CDC 更新覆盖旧文档、按主键幂等max_batch_size 1每条变更立即提交便于在验证阶段逐个观察结果生产环境建议调大以获得更高吞吐源码默认值为 10tls_verify_certificate true和tls_verify_hostname true开启 HTTPS 证书与主机名校验保证传输安全schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST索引不存在时自动创建data_save_mode APPEND_DATA在已有数据基础上追加/更新不清空索引。进一步阅读本教程用到的组件都有独立的专题文档可在仓库中继续深入MySQL-CDCsource 连接器更多启动模式、GTID、表结构变更处理等高级选项Elasticsearchsink 连接器索引模板、动态索引、批量参数、向量字段等完整配置Metadatatransform全部元数据 Key 与 Knowledge Sync 元数据说明Replacetransform正则替换、多字段替换等用法SqltransformSQL 表达式、函数与 UDF 支持。此外仓库中还有 multi-table-cdc.md、mysql-cdc-to-doris.md、mysql-cdc-to-kafka.md 等同类配方可用于对比不同下游的 CDC 同步写法。赞分享数据集成ETL大数据批处理流处理变更数据捕获【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/GitHub_Trending/se/seatunnel点击查看免费下载相关推荐SeaTunnel MySQL CDC 到 Elasticsearch 实战过滤、转换与自定义字段的数据同步配方SeaTunnel MySQL CDC 到 Elasticsearch 实战过滤、转换与自定义字段的数据同步配方 导读 本文是 SeaTunnel 官方 Re数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel 实战PostgreSQL CDC 实时同步到 Iceberg字段整形与 Upsert 完整指南SeaTunnel 实战PostgreSQL CDC 实时同步到 Iceberg字段整形与 Upsert 完整指南 本篇技术指南基于 Apache Sea数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel CDC连接器MySQL到ClickHouse实时同步新范式SeaTunnel CDC连接器MySQL到ClickHouse实时同步新范式 实时数据同步的痛点与破局方案 你是否还在为MySQL到ClickHouse的实数据集成ETL大数据批处理流处理变更数据捕获创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考