完全指南:按管道目标搭建可验证的数据集成作业)
SeaTunnel 场景配方Scenario Recipes完全指南按管道目标搭建可验证的数据集成作业【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本指南面向已经能在本地跑通第一个 SeaTunnel 作业的开发者系统讲解 SeaTunnel 官方场景配方Scenario Recipes体系的选型方法、通用工程要素与实战验证流程。读完本文你将能够根据真实的数据源与目标端组合从 11 个官方配方中快速锁定最接近的管道形态读懂每个配方中env、source、transform、sink四个段落的编写逻辑并掌握插件安装 → 驱动部署 → 前置环境核验 → 作业提交 → 结果断言这条可复现的端到端验证链路。场景配方是 SeaTunnel 官方沉淀的一组最小可验证管道样例覆盖 CDC 实时同步、关系库批式迁移、流式入湖、API 与文件摄取四大类场景。它们与单点功能文档的最大区别在于每个配方都包含具体的前置条件prerequisites、完整配置complete configuration和可预期的结果expected results你可以在自己的环境里照抄、替换、验证而不是面对一堆孤立参数无从下手。配方的推荐阅读时机是你的第一个本地作业已经成功之后——它默认你已经理解了作业的提交方式如./bin/seatunnel.sh --config ./config/xxx.conf -m local因此可以跳过基础铺垫直接进入选型与适配环节。一、按管道目标选择配方官方索引表配方体系的核心是一张目标 → 起点映射表。官方建议不要按顺序通读所有示例而是先判断你的真实管道形态最接近哪一类再从对应配方入手。以下是 recipes 索引页 中的完整选型表管道目标Goal从哪个配方开始Start hereMySQL CDC 写入 Kafka携带元数据 headerMySQL CDC to KafkaMySQL CDC 写入 Elasticsearch带过滤与字段整形MySQL CDC to Elasticsearch关系库之间批式迁移带行级转换JDBC to JDBCJDBC 抽取到对象存储JDBC to S3MySQL 批量抽取为按日期分区的 HDFS ParquetMySQL to HDFSKafka 流式写入 IcebergKafka to IcebergPostgreSQL CDC 写入 IcebergPostgreSQL CDC to IcebergHTTP API 摄取写入 JDBCHTTP to JDBCMySQL CDC 写入 DorisMySQL CDC to Doris文件加载写入 StarRocksFile to StarRocks多表 CDC 编排一作业多表自动路由Multi-Table CDC这 11 个配方在仓库中的完整实现位于 docs/en/getting-started/recipes/ 目录每个配方一个 Markdown 文件。从工程实践角度可以将它们进一步归纳为四条主线CDC 实时同步线MySQL CDC → Kafka / Elasticsearch / DorisPostgreSQL CDC → Iceberg以及多表 CDC → JDBC批式迁移线JDBC → JDBC、JDBC → S3、MySQL → HDFS流式入湖线Kafka → Iceberg数据摄取线HTTP → JDBC、本地文件 → StarRocks。二、如何阅读一个配方四步方法论官方在 overview 中给出了明确的阅读顺序这同样适用于你把配方改造成自己的作业确认源端与目标端组合匹配你的目标管道。例如你的需求是MySQL 变更实时同步到 Kafka就直接进入 MySQL CDC to Kafka 配方而不是从 JDBC to JDBC 开始对比env、source、transform、sink四个段落与你自己的作业。SeaTunnel 的 HOCON 配置天然分为这四层配方里每一层的写法尤其是plugin_output/plugin_input的衔接都是可以复用的骨架改造样例时一次只替换一个系统。先换源端、验证通过后再换目标端避免把源端驱动缺失和目标端 DDL 失败两类问题混在一起排查如果样例依赖 CDC、驱动或额外插件先验证这些前置条件再运行。CDC 类配方对 binlog / 逻辑复制 / 权限有硬性要求驱动 JAR 必须出现在 SeaTunnel 进程的lib目录中这些都属于运行前必须就绪的条件。三、所有配方的通用工程要素在逐类深入之前先归纳 11 个配方共享的工程前置条件——它们是配方能否一次跑通的真正分水岭。3.1 插件安装与plugin_config所有配方都遵循同一套插件安装流程编辑 config/plugin_config只保留配方需要的 connector 条目然后执行安装脚本。例如 JDBC to JDBC 配方要求plugin_config中包含--seatunnel-connectors-- connector-jdbc --end--随后执行cd ${SEATUNNEL_HOME} sh bin/install-plugin.sh ls connectors | rg connector-jdbc仓库根目录的 config/plugin_config 列出了全部可用 connector 的 artifactId 映射如connector-cdc-mysql、connector-jdbc、connector-kafka、connector-iceberg、connector-doris、connector-starrocks、connector-file-local、connector-file-s3、connector-file-hadoop、connector-http-base、connector-elasticsearch等SeaTunnel 依据该文件将配置中的插件名解析为对应的 JAR。注意Sqltransform 内置于发行版中不需要在plugin_config里额外添加条目。3.2 JDBC 驱动必须进入进程的lib目录SeaTunnel Zeta 引擎下所有数据库驱动MySQL Connector/J、PostgreSQL JDBC 驱动等都必须放到${SEATUNNEL_HOME}/lib而不是只存在于你的工作站上。配方中用一条命令验证ls ${SEATUNNEL_HOME}/lib | rg mysql-connector|postgresqlHDFS 场景则要特别注意Zeta 发行版自带 Hadoop JAR先检查lib再决定是否补充依赖不要混用任意版本的 Hadoop 客户端详见 HdfsFile sink 的环境要求。3.3 每个配方都是可断言的配方文档的Validation result / Verified result段落给出了精确的预期结果行数、字段值、目录结构这是它与普通示例的本质区别。验证时建议先记录源端基线如直接执行源查询统计行数再启动作业最后对比目标端。例如 MySQL to HDFS 配方要求先执行分组计数作为基线SELECT DATE(create_time) AS pt_dt, COUNT(*) AS expected_rows FROM trade_db.orders WHERE amount 0 GROUP BY DATE(create_time) ORDER BY pt_dt;四、CDC 实时同步类配方4.1 MySQL CDC → Kafka元数据 header 字段整形这是配方体系中验证最完整的一条官方在 Docker 环境中端到端验证了快照snapshot与 binlog 增量两个阶段。管道形态为MySQL-CDC读取shop.orders→Metadata把库名/表名/行类型暴露为普通字段 →Sql重命名字段并填充自定义字段 →Kafka写 JSON同时把选中的元数据字段搬进 Kafka header。前置条件Docker 验证环境实测MySQL 开启 GTID/binlog配置参考docker/server-gtids/my.cnfMySQL JDBC 驱动 JAR 位于MySQL-CDC插件的lib目录存在 CDC 账号st_user_source具备SELECT、RELOAD、SHOW DATABASES、REPLICATION SLAVE、REPLICATION CLIENT、LOCK TABLES权限Kafka topicrecipe_mysql_orders在作业启动前已创建。完整作业配置Docker 验证通过的原版env { parallelism 1 job.mode STREAMING } source { MySQL-CDC { plugin_output mysql_orders_raw url jdbc:mysql://mysql_cdc_e2e:3306/shop username st_user_source password mysqlpw server-id 5601-5604 table-names [shop.orders] startup.mode initial schema-changes.enabled false } } transform { Metadata { plugin_input mysql_orders_raw plugin_output mysql_orders_with_meta metadata_fields { Database source_database Table source_table RowKind change_type } } Sql { plugin_input mysql_orders_with_meta plugin_output kafka_orders query select id as order_id, order_no, user_id, amount, case when status 0 then CREATED when status 1 then PAID when status 2 then SHIPPED else OTHER end as status_name, source_database, source_table, change_type, CONCAT(source_database, ., source_table) as source_name, mysql_cdc as sync_source from dual where id is not null } } sink { Kafka { plugin_input kafka_orders bootstrap.servers kafkaCluster:9092 topic recipe_mysql_orders format json partition_key_fields [order_id] kafka_headers_fields [source_database, source_table, change_type] } }验证结论Docker E2E 断言项快照阶段为 MySQL 初始行生成了 Kafka 记录记录1001的 payload 包含order_id、status_name、source_name、sync_source快照记录1001的 Kafka headers 为source_databaseshop、source_tableorders、change_typeI且应用kafka_headers_fields后这三个字段从 JSON payload 中移除执行UPDATE shop.orders SET status 2 ... WHERE id 1001后1001最新记录为change_typeU、status_nameSHIPPED插入1003后其最新记录为change_typeI、status_nameCREATED。kafka_headers_fields的源码级行为Kafka sink 对 header 字段有一组强校验见 KafkaSinkWriter.javakafka_headers_fields不支持NATIVE格式使用 JSON、TEXT 等其他格式header 字段不能与partition_key_fields重叠否则抛 Field %s cannot be in both partition_key_fields and kafka_headers_fieldsheader 字段不能与kafka_message_value_fields重叠每个 header 字段必须真实存在于行类型rowType中否则抛 Header field not found见同文件 L321-L339。4.2 MySQL CDC → Elasticsearch过滤 字段清洗 自定义字段该配方持续把客户档案从 MySQL 同步到 ElasticsearchMySQL-CDC读取crm.customer_profile→Metadata暴露库/表/行类型 →Replace去掉电话号码中的-→Sql过滤行并派生status_name与sync_source→Elasticsearch按主键写文档。关键点Elasticsearch 8.9.0 要求 HTTPS 认证需把 CA 证书复制到每个 SeaTunnel 节点并用tls_truststore_path引用本地路径每个并发运行的 MySQL CDC 作业必须使用唯一的server-id范围。核心配置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 } }为什么这样设计Replace只做轻量字段归一化去-Sql独占业务过滤逻辑——按不可变主键id 1000过滤、把status映射为ACTIVE/FROZEN、填充常量sync_source。配方特别提醒不要用is_deleted、status这类可变字段做过滤否则当已索引行不再匹配过滤条件时upsert 型 sink 不会为旧文档发出删除事件软删除场景应把删除标志保留在 Elasticsearch 中并在查询时过滤或使用能显式把状态迁移转换为删除事件的管道。4.3 MySQL CDC → Doris持续更新 删除回放配方目标捕获 MySQL 行级变更并持续更新 Doris 表。前提包括MySQL binlog 就绪log_bin ON、binlog_format ROW、binlog_row_image FULL否则在my.cnf中配置后重启、CDC 用户具备复制权限、源表有稳定主键。核心配置env { parallelism 1 job.mode STREAMING checkpoint.interval 5000 } source { MySQL-CDC { plugin_output orders_cdc parallelism 1 startup.mode initial server-id 5652 username st_user_source password mysqlpw table-names [inventory.orders] url jdbc:mysql://mysql:3306/inventory } } sink { Doris { plugin_input orders_cdc fenodes doris-fe:8030 username root password database sync_demo table orders sink.label-prefix orders-cdc sink.enable-delete true schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST doris.config { format csv column_separator , } } }验证时先在 MySQL 上执行INSERT 1003、UPDATE 1001 → PAID、DELETE 1002再查询 Doris 确认最终状态。常见坑包括server-id与其它 MySQL 副本或 CDC 作业冲突sink.label-prefix被多个运行中作业复用导致 Doris stream load 冲突开启删除回放但 Doris 表模型不支持预期删除行为源表无稳定主键导致自动建表与下游 upsert 不确定。4.4 PostgreSQL CDC → Iceberg字段整形 upsert该配方通过 PostgreSQL 逻辑复制pgoutput捕获sales.inventory.customer_orders的快照与后续变更snapshot / update / insert / delete在 Iceberg 中维护当前状态。前置条件主库开启wal_level logical并重启创建带REPLICATION属性的只读 CDC 角色并授予CONNECT、schemaUSAGE、表SELECT在pg_hba.conf中为逻辑复制连接添加规则注意规则指向实际数据库sales而非物理复制用的replication关键字由管理员创建 publication使 CDC 用户保持只读。源表需设置REPLICA IDENTITY FULL。完整配置env { parallelism 1 job.mode STREAMING checkpoint.interval 3000 } source { Postgres-CDC { plugin_output postgres_orders_raw url jdbc:postgresql://postgresql.example.com:5432/sales username seatunnel_cdc password change_me database-names [sales] schema-names [inventory] table-names [sales.inventory.customer_orders] startup.mode initial decoding.plugin.name pgoutput slot.name ${slot_name} debezium { publication.name seatunnel_sales_orders_pub publication.autocreate.mode disabled } } } transform { Sql { plugin_input postgres_orders_raw plugin_output iceberg_customer_orders query select id, trim(customer_name) as customer_name, amount, upper(status) as status_name, updated_at, postgresql_cdc as sync_source from dual } } sink { Iceberg { plugin_input iceberg_customer_orders catalog_name recipe_catalog iceberg.catalog.config { type hadoop warehouse ${warehouse} } namespace sales_analytics table customer_orders iceberg.table.primary-keys id iceberg.table.upsert-mode-enabled true schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST data_save_mode APPEND_DATA } }运行前需把${slot_name}与${warehouse}替换为字面量例如slot.name seatunnel_sales_orders、warehouse file:///tmp/seatunnel/iceberg/postgres-cdc-recipe/。slot 的源码级约束slot.name是 PostgresIncrementalSourceOptions.java 定义的核心参数同一 PostgreSQL 实例上的每个并发 CDC 作业必须使用不同的 slot 名PostgresSourceFetchTaskContext.java 中明确提示为同一数据库主机设置多个 connector 时请确保每个使用独立的复制槽名称。运维上要监控pg_replication_slots停止的消费者可能无限期保留 WAL只有作业永久下线后才应删除 slot。4.5 Multi-Table CDC一作业多表自动路由该配方用一个作业捕获多张上游表并把每张表自动路由到各自的下游表。MySQL-CDC通过database-patterntable-pattern一次订阅inventory库的orders/customers/products三张表Jdbcsink 借助占位符${table_name}、${primary_key}自动生成st_上游表名系列目标表依赖generate_sink_sql true与 PostgreSQL 用户对publicschema 的CREATE权限。source { MySQL-CDC { plugin_output mysql_multi startup.mode initial server-id 5652 username st_user_source password mysqlpw database-pattern inventory table-pattern inventory\\.(orders|customers|products) url jdbc:mysql://mysql:3306/inventory } } sink { Jdbc { plugin_input mysql_multi driver org.postgresql.Driver url jdbc:postgresql://postgresql:5432/sync_demo username st_user_sink password pgpw generate_sink_sql true database sync_demo table public.st_${table_name} primary_keys [${primary_key}] schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST data_save_mode APPEND_DATA } }验证方法先查information_schema.tables确认自动创建了多个st_%表再在每张上游表上分别插入/更新最后逐表查询确认每张下游表只收到自己的源数据。高频坑位table-pattern的正则在 HOCON 中未正确转义字面点号需要用\\.未配置占位符路由导致多张源表被误写入同一张目标表下游 schema 与表占位符的用法在不同数据库中语义不同。更完整的架构说明见 Multi-table 同步架构。五、批式迁移类配方5.1 JDBC → JDBC行级过滤与整形适用场景关系库间批式迁移写前需要对行做过滤或重塑。示例从 MySQL 读订单仅保留已支付订单status PAID、姓名转大写、金额精度归一化、填充source_system常量字段写入 PostgreSQL。运行方式cd ${SEATUNNEL_HOME} ./bin/seatunnel.sh --config ./config/jdbc-to-jdbc.conf -m local完整配置env { parallelism 1 job.mode BATCH } source { Jdbc { plugin_output mysql_orders driver com.mysql.cj.jdbc.Driver url jdbc:mysql://mysql.example.com:3306/source_db?useSSLfalseallowPublicKeyRetrievaltrue username root password password query SELECT id, customer_name, amount, status FROM orders } } transform { Sql { plugin_input mysql_orders plugin_output paid_orders query SELECT id, UPPER(customer_name) AS customer_name, CAST(amount AS DECIMAL(12, 2)) AS amount, MYSQL AS source_system FROM dual WHERE status PAID } } sink { Jdbc { plugin_input paid_orders driver org.postgresql.Driver url jdbc:postgresql://postgresql.example.com:5432/target_db username test password test query INSERT INTO public.paid_orders (id, customer_name, amount, source_system) VALUES (?, ?, ?, ?) } }关键约定dual是 SeaTunnel 默认 SQL transform 引擎使用的虚拟输入表不是MySQL 或 PostgreSQL 中的真实表plugin_output/plugin_input串联三个阶段转换后的列顺序必须与 sinkquery中的四个占位符严格对齐。验证结果1001为ALICE CHEN / 120.50 / MYSQL1003为CAROL WU / 42.00 / MYSQL1002因未支付被过滤。常见坑包括驱动版本不兼容、示例主机名未替换、字段顺序与 INSERT 列不一致、目标表非空导致第二次运行主键冲突教程场景先 TRUNCATE生产请选择 upsert 策略、源 DECIMAL 超过目标DECIMAL(12, 2)精度、批式运行期间源数据发生变化需持续捕获变更时应改用 CDC 源。5.2 JDBC → S3查询结果导出对象存储该配方把 MySQL 查询结果以 JSON Lines 导出到 S3。S3 连接需要hadoop-aws与aws-java-sdk-bundle两个 JAR 位于${SEATUNNEL_HOME}/lib。核心配置要点sink { S3File { plugin_input orders_jdbc bucket s3a://company-data-lake path /seatunnel/orders/ fs.s3a.endpoint s3.us-east-1.amazonaws.com fs.s3a.aws.credentials.provider org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider access_key your-access-key secret_key your-secret-key file_format_type json row_delimiter \n custom_filename true file_name_expression orders filename_extension json single_file_mode true is_enable_transaction false schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST data_save_mode APPEND_DATA } }验证用aws s3 ls/aws s3 cp检查目标前缀下的对象与内容。坑位驱动只在工作站而没进libbucket与path混用bucket 放桶名、path 放前缀凭据提供者与认证方式不匹配大表一条无界查询导出无过滤无分区固定文件名只适合单文件教程场景若重新启用事务需在file_name_expression中保留${transactionId}S3 兼容端点却仍指向 AWS。5.3 MySQL → HDFS按日期分区的 Snappy Parquet配方目标把 MySQL 订单批量加载为 HDFS 上按日期分区的 Snappy 压缩 Parquet属于批式快照而非 CDC。Sqltransform 把id重命名为order_id、状态转大写、过滤负金额、派生pt_dt分区键。核心配置env { job.name mysql_to_hdfs_batch_dw job.mode BATCH parallelism 4 } source { Jdbc { plugin_output src_mysql_orders url jdbc:mysql://localhost:3306/trade_db?useSSLfalseserverTimezoneUTCrewriteBatchedStatementstrue driver com.mysql.cj.jdbc.Driver username test_user password test_password table_path trade_db.orders query select id, order_no, user_id, amount, status, create_time, date(create_time) as create_date from trade_db.orders partition_column id partition_num 4 partition_lower_bound 1 partition_upper_bound 10000000 fetch_size 2000 } } transform { Sql { plugin_input src_mysql_orders plugin_output dwd_orders query select id as order_id, order_no, user_id, amount, upper(status) as order_status, create_time, FORMATDATETIME(create_date, yyyy-MM-dd) as pt_dt from src_mysql_orders where amount 0 } } sink { HdfsFile { plugin_input dwd_orders fs.defaultFS hdfs://namenode:8020 path /user/hive/warehouse/dwd.db/dwd_orders_df file_format_type parquet partition_by [pt_dt] compress_codec snappy schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST data_save_mode APPEND_DATA } }这里有两个值得注意的设计date(create_time)在 MySQL 侧求值FORMATDATETIME(create_date, yyyy-MM-dd)在 SeaTunnel SQL 侧生成分区键partition_column/partition_num/partition_lower_bound/partition_upper_bound演示的是 JDBC 分片读split-read配置对 5 行数据而言并非性能调优扩量时需按真实数据调整。验证要点hdfs dfs -ls -R应看到pt_dt2026-08-23/与pt_dt2026-08-24/两个分区目录用 Parquet 读取器核对 4 行记录负金额的 order 5 被过滤检查元数据确认 Snappy 压缩不要用cat读二进制Hive 风格目录不会自动建/注册 Hive 表目录存在不等于完整验证。该作业使用APPEND_DATA重复运行会产生重复记录换用新的输出路径做下一轮验证。六、流式入湖类配方Kafka → Iceberg适用场景把 Kafka 流式事件落入 Iceberg 表供下游分析。管道形态KafkasourceJSON schema 声明→IcebergsinkHadoop catalog upsert schema evolution。env { parallelism 1 job.mode STREAMING checkpoint.interval 5000 } source { Kafka { plugin_output orders_kafka topic orders bootstrap.servers kafka:9092 consumer.group seatunnel-orders start_mode earliest format json schema { fields { id bigint customer_id bigint total_amount decimal(16, 2) event_date string } } } } sink { Iceberg { plugin_input orders_kafka catalog_name seatunnel_demo namespace lakehouse table orders iceberg.catalog.config { type hadoop warehouse file:///tmp/seatunnel/iceberg/warehouse-demo } iceberg.table.primary-keys id iceberg.table.partition-keys event_date iceberg.table.upsert-mode-enabled true iceberg.table.schema-evolution-enabled true case_sensitive true } }验证方式先ls /tmp/seatunnel/iceberg/warehouse-demo/lakehouse/orders确认元数据与数据文件生成再用 Spark/Trino 等 Iceberg 兼容引擎查询两条样例消息应得到COUNT(*) 2。常见坑Kafka 中的 JSON 与 source schema 不匹配流式管道禁用 checkpoint 会削弱重启与一致性warehouse 路径对引擎进程不可写记录没有稳定主键却开启了 upsert 模式。Flink/Spark 运行时还需要按环境补充hive-exec、libfb303等 Iceberg 依赖。七、数据摄取类配方7.1 HTTP → JDBCAPI 摄取自动建表适用场景从 HTTP API 拉取结构化数据存入关系库。前置要点先curl检查响应体结构样例端点返回顶层字段c_string、c_int若真实 API 把记录嵌套在其它字段下需先配置json_field或content_fieldsink 用户需要对publicschema 有USAGE, CREATE权限因为配方使用generate_sink_sql true自动建表。source { Http { plugin_output http_orders url http://mockserver:1080/example/http method GET format json schema { fields { c_string string c_int int } } } } sink { Jdbc { plugin_input http_orders driver org.postgresql.Driver url jdbc:postgresql://postgresql:5432/test?loggerLevelOFF username test password test generate_sink_sql true database test table public.http_orders primary_keys [c_string] batch_size 100 } }验证确认无 HTTP 解析/JDBC DDL 错误对比目标表行数与 API 响应。坑位schema 与实际字段名/类型不符嵌套数据未配置content_field/json_field源 API 有分页或限流却被当作单页端点自动建表时选的主键不能唯一标识记录。7.2 本地文件 → StarRocksCSV 批量导入适用场景把本地 CSV/文本文件导入 StarRocks 做分析查询。前置connector-file-local与connector-starrocks两个插件、lib中放 MySQL JDBC 驱动StarRocks sink 依赖、预创建 StarRocks 主键表该配方用schema_save_mode IGNORE因为本地文件源不提供主键元数据给 StarRocks 自动 DDL。source { LocalFile { plugin_output customers_file path /tmp/seatunnel/input/customers.csv file_format_type csv csv_use_header_line true schema { fields { id bigint name string city string updated_at timestamp } } } } sink { StarRocks { plugin_input customers_file nodeUrls [starrocks-fe:8030] base-url jdbc:mysql://starrocks-fe:9030/sync_demo username root password database sync_demo table customers batch_max_rows 1000 schema_save_mode IGNORE starrocks.config { format JSON strip_outer_array true } } }坑位配了nodeUrls却漏了base-url文件有表头但未设csv_use_header_line true源 schema 与文件分隔符/时间戳格式不匹配运行前未建目标表。八、从配方到生产跨场景注意事项汇总纵览 11 个配方可以提炼出一套可复用的生产化检查清单CDC 类作业的 binlog/逻辑复制是硬前提MySQL 需log_bin ON、binlog_format ROW、binlog_row_image FULL检查命令见各配方SHOW VARIABLES WHERE variable_name IN (log_bin, binlog_format, binlog_row_image)PostgreSQL 需wal_level logical且修改后重启。CDC 账号权限、server-id唯一性每个并发作业一个独立值或范围、PostgreSQLslot.name唯一性缺一不可。这些参数在 MySqlIncrementalSourceOptions.java 与 PostgresIncrementalSourceOptions.java 中均有对应定义源码中还包含了对冲突场景的显式告警。插件与驱动按版本对齐config/plugin_config中的 connector 条目、发行版版本、lib中的 JDBC 驱动三者必须匹配Hadoop 相关依赖先检查lib再补充避免混入任意版本。schema 与 data 保存模式决定幂等性多数配方用schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST与data_save_mode APPEND_DATA。前者让作业自动建表依赖源端提供主键元数据后者意味着重复运行会追加数据——教程场景用新输出路径或 TRUNCATE生产场景必须换成 upsert / 事务策略。流式作业依赖 checkpointjob.mode STREAMING时配合checkpoint.interval配方中常用 30005000msIceberg 等目标端的变更在 checkpoint 成功提交后才可见禁用 checkpoint 会削弱重启恢复与一致性。过滤与 upsert 的陷阱CDC 管道中过滤条件应基于不可变主键而非可变业务字段启用 upsert 模式必须有明确的primary-keys且源数据要有稳定主键。先验证、再断言、最后下结论每个配方都给出了精确的预期结果表如 JDBC→JDBC 的两行结果、PG CDC→Iceberg 的最终主键集、MySQL→HDFS 的分区目录树这是判断管道是否真正工作的客观标准。九、延伸阅读配方总入口Scenario Recipes overview前置技能Run your first job、插件部署与下载核心连接器文档MySQL-CDC source、PostgreSQL CDC source、Kafka sink、JDBC source、JDBC sink、Elasticsearch sink、Doris sink、Iceberg sink、S3File sink、HdfsFile sink、Http source、StarRocks sink、LocalFile source相关 transform 文档Metadata transform、Sql transform、Replace transform、SQL functions架构参考Multi-table 同步架构【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考