SeaTunnel核心原理

发布时间:2026/9/30 19:27:45
SeaTunnel核心原理 1 背景目标基于SeaTunnel构建数据集成平台数据集成主要用于数据采集和推送打通不同数据源之间的数据流转和汇聚实现数据统一管理、流转和共享。从维护、治理、效能三方面分析痛点挑战①高可靠能够在各种故障情况下保持数据不丢失保障数据一致性和任务的稳定运行②高吞吐能够处理大规模数据流实现高并发和大批量数据传输③低延迟满足实时数据处理和快速响应的业务需求④易维护简洁的配置和自动化的监控减轻运维负担便于快速发现和解决故障确保系统长期可用核心目标提升效率同步效率提升3倍 | 开发标准规范覆盖率 90%↑ | 赋能其他团队2 SeaTunnel架构图3 SeaTunnel核心原理3.1 连接器插件架构SeaTunnel 连接器基于统一插件化架构分为SourceAPI和Sink API两大部分通过 Enumerator 分发分片、Task Group 并行执行读写、Committer 聚合提交状态。组件所属职责EnumeratorSource API读取 Catalog生成 Split 列表并分发给 ReaderReaderSource API接收 Split读取数据并发送 row 给 WriterState APISource Sink保存/恢复快照状态支持断点续传WriterSink API接收数据行写入目标数据源Committer APISink API聚合所有 Writer 的提交状态CommitterSink API两阶段提交保证 Exactly-Once 语义3.2 分布式快照与 CheckpointSeaTunnel 实现统一 Checkpoint 抽象层屏蔽底层引擎差异Zeta / Flink / Spark上层连接器只需实现一套 API。底层引擎Zeta/Flink/Spark各自有 Checkpoint 实现但最终都收敛到统一的 SeaTunnel Checkpoint 流程保证跨引擎一致性语义。3.3 类型映射与自动建表Source 推送表结构convert外部数据库类型 → SeaTunnel ColumnColumn convert(T typeDefine);Sink 写表结构reconvert是把SeaTunnel Column → 外部数据库类型T reconvert(Column column);中间层是SeaTunnel 类型涉及到类型提升。SeaTunnel 内部类型中间类型桥梁SeaTunnelDataType.javaSeaTunnelDataTypeT ├── BasicType (STRING, BOOLEAN, BYTE, SHORT, INT, LONG, FLOAT, DOUBLE, VOID) ├── DecimalType (precision, scale) ├── LocalTimeType (DATE, TIME, TIMESTAMP, TIMESTAMP_TZ) ├── PrimitiveByteArrayType (BYTES) ├── ArrayTypeT,E ├── MapTypeK,V └── SeaTunnelRowType (嵌套行)TypeConverter.java public interface TypeConverterT { // 外部 DB 类型 → SeaTunnel Column Column convert(T typeDefine); // SeaTunnel Column → 外部 DB 类型 T reconvert(Column column); }自动建表机制 (SaveMode)建表流程DefaultSaveModeHandler → Catalog.createTable() → CreateTableSqlBuilder → TypeConverter.reconvert() → DDL核心接口 // Sink 连接器声明支持 SaveMode public interface SupportSaveMode { OptionalSaveModeHandler getSaveModeHandler(); } // SaveMode 处理器 public interface SaveModeHandler extends AutoCloseable { void handleSchemaSaveMode(); // 处理表结构 void handleDataSaveMode(); // 处理数据 }SaveMode类型行为RECREATE_SCHEMASchemaSaveMode(表结构处理策略)表不存在则创建表存在则删除重建CREATE_SCHEMA_WHEN_NOT_EXISTSchemaSaveMode(表结构处理策略)表不存在则创建表存在则跳过ERROR_WHEN_SCHEMA_NOT_EXISTSchemaSaveMode(表结构处理策略)表不存在则报错IGNORESchemaSaveMode(表结构处理策略)不处理DROP_DATADataSaveMode数据处理策略保留表结构清空数据 (TRUNCATE)APPEND_DATADataSaveMode数据处理策略保留表结构追加写入CUSTOM_PROCESSINGDataSaveMode数据处理策略用户自定义处理逻辑ERROR_WHEN_DATA_EXISTSDataSaveMode数据处理策略目标有数据则报错3.4 多表同步3.4.1 多表同步架构Source 多表读取Sink 多表写入SplitEnumerator 按表逐一生成 Split单表内按主键范围拆分为多个 BlockSplit Hash路由分发给多个 Reader 并行读取Reader 循环读取分配到的 Split按表迭代生成 Sub Sink拆分独立的 Writer 和 Committer数据按表路由接收各表独立提交共享连接资源降低连接开销3.4.2 Source处理多表读取1Source 端逐表读取实现分析Source 端实现逐表读取的核心设计是Split 携带表标识Enumerator 逐表生成 SplitReader 读取时自动设置行的 tableId。整个过程涉及三个组件的协作。①Enumerator (Master) pendingTables: [table_A, table_B, table_C] ← 初始化时从配置获取所有表 run() { while (!pendingTables.isEmpty()) { tablePath pendingTables.poll(); // ① 逐表取出 splits splitter.generateSplits(table); // ② 为该表生成 splits addPendingSplit(splits); // ③ hash 分配给各 Reader assignSplit(readers); // ④ 下发 } signalNoMoreSplits(); // ⑤ 所有表处理完毕 } ②Reader (Worker)【split 下发】 splits deque: [split_A_0, split_A_1, split_B_0, split_C_0, ...] pollNext() { split splits.poll(); // ⑥ 逐个取 split inputFormat.open(split); // ⑦ 根据 split.tablePath 定位表 while (!reachedEnd()) { row inputFormat.nextRecord(); // ⑧ 读取一行 // row.tableId 已设置为 split.tablePath output.collect(row); // ⑨ 发出数据 } inputFormat.close(); }关键代码解析① Enumerator 逐表迭代pendingTables 是一个队列逐个 poll() 取出表每张表独立生成 split 后立即分配。这实现了逐表切分的语义。② Split 携带表标识每个 Split 都绑定了 tablePath。当 Reader 拿到一个 split 时就知道该从哪张表读取数据。③ Reader 逐 Split 处理自然实现逐表Reader 的 splits 队列中可能混合了不同表的 split因为 Enumerator 按 hash 分配。Reader 不关心顺序逐个处理即可每个 split 的 tablePath 决定了读取哪张表。④ InputFormat 设置行的 tableId这是逐表读取的关键一环open(split) 时从 split 的 tablePath 获取表的 schema 和标识nextRecord() 产出的每一行都通过 setTableId() 打上所属表的标记。① Enumerator 逐表迭代 // JdbcSourceSplitEnumerator.java private final ConcurrentLinkedQueueTablePath pendingTables; // 待处理的表队列 Override public void run() throws Exception { SetInteger readers context.registeredReaders(); while (!pendingTables.isEmpty()) { synchronized (stateLock) { TablePath tablePath pendingTables.poll(); // 逐表取出 LOG.info(Splitting table {}., tablePath); CollectionJdbcSourceSplit splits splitter.generateSplits(tables.get(tablePath)); LOG.info(Split table {} into {} splits., tablePath, splits.size()); addPendingSplit(splits); // 分配给 Reader } synchronized (stateLock) { assignSplit(readers); // 下发 split } } splitter.close(); readers.forEach(context::signalNoMoreSplits); // 所有表都处理完了 } ② Split 携带表标识 // JdbcSourceSplit.java public class JdbcSourceSplit implements SourceSplit { private final TablePath tablePath; // ← 每个 split 知道自己属于哪张表 private final String splitId; private final String splitQuery; private final String splitKeyName; private final SeaTunnelDataType splitKeyType; private final Object splitStart; private final Object splitEnd; } ③ Reader 逐 Split 处理自然实现逐表 // JdbcSourceReader.java private final DequeJdbcSourceSplit splits new ConcurrentLinkedDeque(); Override public void pollNext(CollectorSeaTunnelRow output) throws Exception { synchronized (output.getCheckpointLock()) { JdbcSourceSplit split splits.poll(); // 逐个取 split if (null ! split) { try { inputFormat.open(split); // 打开该 split 对应的表查询 while (!inputFormat.reachedEnd()) { SeaTunnelRow row inputFormat.nextRecord(); output.collect(row); // 逐行输出 } } finally { inputFormat.close(); // 关闭当前 split } } else if (noMoreSplit splits.isEmpty()) { context.signalNoMoreElement(); // 全部读完 } else { Thread.sleep(1000L); // 等待新 split 分配 } } } ④ InputFormat 设置行的 tableId // JdbcInputFormat.java public void open(JdbcSourceSplit inputSplit) throws IOException { // 从 split 获取表信息 splitTableSchema tables.get(inputSplit.getTablePath()).getTableSchema(); splitTableId inputSplit.getTablePath().toString(); // ← 记住当前表 ID // 生成该 split 的查询语句并执行 statement chunkSplitter.generateSplitStatement(inputSplit, splitTableSchema); resultSet statement.executeQuery(); hasNext resultSet.next(); } public SeaTunnelRow nextRecord() { SeaTunnelRow seaTunnelRow jdbcRowConverter.toInternal(resultSet, splitTableSchema); seaTunnelRow.setTableId(splitTableId); // ← 每行数据标记所属表 seaTunnelRow.setRowKind(RowKind.INSERT); hasNext resultSet.next(); return seaTunnelRow; }2数据行如何标识所属表Source 在产出每一行数据时会设置 row.setTableId(tablePath.toString())。这个 tableId 会跟随数据行在整个 Pipeline 中流动直到 Sink 端用它来决定写入哪张表。数据行如何标识所属表 — SeaTunnelRow.tableIdpublic class SeaTunnelRow { private String tableId ; // ← 表标识多表路由的关键 private RowKind rowKind RowKind.INSERT; // INSERT/UPDATE/DELETE private final Object[] fields; }3Source 端多表支持getProducedCatalogTables() 返回多个 CatalogTable引擎据此构建多表 DAG。source { Jdbc { url jdbc:mysql://localhost:3306/mydb driver com.mysql.cj.jdbc.Driver # 方式一table_list 显式列出多张表 table_list [ { table_path mydb.users }, { table_path mydb.orders, partition_column id }, { table_path mydb.products } ] # 方式二table_path 正则 (use_regextrue) # table_path mydb.order_\\d # use_regex true } } // 内部实现 // JdbcSourceTableConfig 解析 table_list 配置 public static ListJdbcSourceTableConfig of(ReadonlyConfig config) { if (config.getOptional(TABLE_LIST).isPresent()) { return config.get(TABLE_LIST); // 多表配置 } else { return singletonList(/*单表配置*/); } } // JdbcSource 构造时获取所有表的元数据 this.jdbcSourceTables JdbcCatalogUtils.getTables(connectionConfig, tableConfigList); // 返回 MapTablePath, JdbcSourceTable // Source 声明产出多张表 Override public ListCatalogTable getProducedCatalogTables() { return jdbcSourceTables.values().stream() .map(JdbcSourceTable::getCatalogTable) .collect(Collectors.toList()); }4DAG解析模式 — ParsingModeParsingMode.javapublic enum ParsingMode { SINGLENESS, // 每张表独立 Source Sink最简单 SHARDING, // 多个分片表合并为一个 Source 一个 Sink分库分表场景 MULTIPLEX // 一个 Source 产出多表各表使用独立 Sink已废弃 }模式DAG结构使用场景SINGLENESSsource(A)→sink(A), source(B)→sink(B)表之间完全无关SHARDINGsource(A1,A2,A3)→sink(A)分库分表合并到一张目标表MULTIPLEXsource(A,B,C)→sink(A), →sink(B), →sink(C)多表同步 (默认行为)在 SHARDING 模式下FactoryUtil 只保留第一张 CatalogTable 的结构if (options.get(DAG_PARSING_MODE) ParsingMode.SHARDING) { CatalogTable catalogTable catalogTables.get(0); catalogTables.clear(); catalogTables.add(catalogTable); // 多个分片表视为同一结构 }5引擎DAG构建流程 (MultipleTableJobConfigParser)--Zeta 引擎parse() ├── 1. parseSource() → 生成 SourceAction ListCatalogTable │ └── 每个 CatalogTable 对应一个 Tuple2CatalogTable, Action │ └── 存入 tableWithActionMap[tableId] ListTuple2CatalogTable, Action │ ├── 2. parseTransforms() → Transform 处理多表输入 │ └── 3. parseSink() → 为每张表创建 per-table sink ├── 遍历 inputVertices 中的每个 CatalogTable ├── 为每张表调用 FactoryUtil.createAndPrepareSink() ├── 生成多个 SinkAction └── tryGenerateMultiTableSink() ├── 检查所有 sink 是否都实现 SupportMultiTableSink ├── 如果是 → 合并为一个 MultiTableSink └── 如果否 → 保持独立的 SinkAction (sink template 模式)关键代码 — 尝试合并为 MultiTableSinkprivate OptionalSinkAction?,?,?,? tryGenerateMultiTableSink( ListSinkAction?,?,?,? sinkActions, ...) { // 条件所有 sink 都必须声明支持多表 if (sinkActions.stream().anyMatch( action - !(action.getSink() instanceof SupportMultiTableSink))) { log.info(Unsupported multi table sink api, rollback to sink template); return Optional.empty(); } // 收集所有 per-table sink MapTablePath, SeaTunnelSink sinks new HashMap(); sinkActions.forEach(action - { sinks.put(action.getConfig().getTablePath(), action.getSink()); }); // 创建 MultiTableSink 包装 SeaTunnelSink?,?,?,? sink FactoryUtil.createMultiTableSink(sinks, options, classLoader); return Optional.of(new SinkAction(..., sink, ...)); }3.4.3 Sink处理多表写入Sink处理多表写入按表迭代生成Sub Sink拆分 Writer/Committer数据按表路由接收共享链接资源1MultiTableSink — 多表 Sink 包装器每个 subtask 为每张表创建 replicaNum 个 Writer 副本用于提高写入吞吐。核心结构 public class MultiTableSink implements SeaTunnelSink... { private final MapTablePath, SeaTunnelSink sinks; // 每表一个独立 Sink private final int replicaNum; // 每表的并行写入副本数 } Write创建 Override public SinkWriterSeaTunnelRow, MultiTableCommitInfo, MultiTableState createWriter( SinkWriter.Context context) { MapSinkIdentifier, SinkWriterSeaTunnelRow, ?, ? writers new HashMap(); for (int i 0; i replicaNum; i) { for (TablePath tablePath : sinks.keySet()) { int index context.getIndexOfSubtask() * replicaNum i; writers.put( SinkIdentifier.of(tablePath.toString(), index), sinks.get(tablePath).createWriter(new SinkContextProxy(index, ...)) ); } } return new MultiTableSinkWriter(writers, replicaNum, sinkWritersContext); }2MultiTableSinkWriter — 多表路由分发路由策略:有主键: hash(pk_value) % queueSize无主键: random() % queueSizeMultiTableSinkWriterQueue[0] ← WriterRunnable[0] table_A Writer table_C WriterQueue[1] ← WriterRunnable[1] table_B Writer table_D Writer写入路由逻辑Override public void write(SeaTunnelRow element) throws IOException { OptionalInteger primaryKey sinkPrimaryKeys.get(element.getTableId()); if (primaryKey ! null primaryKey.isPresent()) { // 有主键 → 按主键值 hash 路由保证相同 key 到同一 queue有序性 Object pkValue element.getField(primaryKey.get()); int index Math.abs(pkValue.hashCode()) % blockingQueues.size(); blockingQueues.get(index).offer(element); } else { // 无主键 → 随机路由负载均衡 int index random.nextInt(blockingQueues.size()); blockingQueues.get(index).offer(element); } }3MultiTableWriterRunnable — 异步消费线程public class MultiTableWriterRunnable implements Runnable { private final MapString, SinkWriterSeaTunnelRow, ?, ? tableIdWriterMap; private final BlockingQueueSeaTunnelRow queue; Override public void run() { while (true) { SeaTunnelRow row queue.poll(100, TimeUnit.MILLISECONDS); if (row null) continue; // 按 tableId 路由到对应的 per-table Writer SinkWriterSeaTunnelRow, ?, ? writer tableIdWriterMap.get(row.getTableId()); synchronized (this) { writer.write(row); } } } }4资源共享 — SupportResourceShare效果同一 subtask 上的所有表 Writer 共享一个连接池而非每表一个。多表场景下如果每张表都创建独立的数据库连接池资源开销巨大。SupportResourceShare 允许共享连接池SupportResourceShare.javapublic interface SupportResourceShareT { // 创建共享资源如连接池只调用一次 MultiTableResourceManagerT initMultiTableResourceManager(int tableSize, int queueSize); // 注入共享资源到每个 per-table Writer void setMultiTableResourceManager(MultiTableResourceManagerT manager, int queueIndex); }举个JDBC链接的例子// JdbcSinkWriter.initMultiTableResourceManager public MultiTableResourceManagerConnectionPoolManager initMultiTableResourceManager( int tableSize, int queueSize) { HikariDataSource ds new HikariDataSource(); ds.setMaximumPoolSize(queueSize); // 连接池大小 queue 数 ds.setJdbcUrl(url); // ... 配置连接池 return new JdbcMultiTableResourceManager(new ConnectionPoolManager(ds)); }5Schema Evolution 支持MultiTableSinkWriter 还支持多表场景下的 Schema EvolutionDDL 变更同步。DDL 变更会精准路由到对应表的 Writer同时通过 synchronized 确保与数据写入互斥。Override public void applySchemaChange(SchemaChangeEvent event) { // 按 tablePath 找到对应的 per-table writer for (entry : sinkWritersWithIndex) { if (entry.getKey().getTableIdentifier().equals(event.tablePath().getFullName())) { synchronized (runnable.get(queueIndex)) { // 暂停该 queue 的写入 ((SupportSchemaEvolutionSinkWriter) entry.getValue()) .applySchemaChange(event); // 执行 DDL } } } }3.5 一读多写共享读取缓存Source 仅读取一次数据同时写入多个目标MQ-Sink / DB-Sink / File-Sink。一读多写共享读取缓存一次写入多种目标Source仅读取一次写入速率会相互影响注意各 Sink 写入速率会相互影响最慢的 Sink 会成为整体瓶颈需合理规划各目标写入能力。3.6 CDC 功能详解3.6.1 基本流程分成两个阶段①快照读取阶段用于读取表的历史数据快照Snapshot全量、回填Backfill②增量读取阶段用于读取表的增量日志更改数据用MySQL-CDC举例说一下基本原理数据来源于数据库的二进制日志Binlog。CDC连接器底层基于Debezium会将自己伪装成一个MySQL副本持续接收并解析原始的行级别变更事件流。3.6.2 Exactly-Once 保障1快照阶段读不丢不重数据库持续发生更改如何保证读取过程数据一致性不丢失、不重复流程步骤快照阶段读不丢不重读 split 前记录low-watermark读取 Split[Start~End] 数据读 split 后记录high-watermark若 high low 则执行 Merge上报 low/high-watermarkMySQL内存表缓存整个 Split 数据过滤 low/high 之间的 binlog保证不丢不重。Oracle可原子化操作select ... for scn直接输出精确数据。举个例子老师让数数班级同学同学比较调皮跑来跑去如何做SeaTunnel 是这么做的先拍照快照先把班里的人数一遍记在小本本上内存缓存先不告诉校长下游。看监控回填去调取数数这段时间的监控录像Binlog 日志。修正记录合并如果监控显示有人刚进来但你没数到 - 加上。如果监控显示有人刚跑了但你数进去了 - 划掉。如果监控显示有人换了件衣服 - 把记录改成新衣服。源码分析第一步生成“看监控”的任务当快照读完后SeaTunnel 会自动生成一个任务专门去读 “开始数数”到“数数结束” 这段时间的日志。MySqlSnapshotFetchTask.javaprivate IncrementalSplit createBackfillBinlogSplit( SnapshotSplitChangeEventSourceContext sourceContext) { return new IncrementalSplit( split.splitId(), Collections.singletonList(split.getTableId()), sourceContext.getLowWatermark(), // 监控录像开始时间低水位线 sourceContext.getHighWatermark(), // 监控录像结束时间高水位线 new ArrayList()); }第二步看着监控改作业核心逻辑这是最关键的一步。SeaTunnel 会拿着日志监控录像去修改内存里的快照数据小本本。JdbcSourceFetchTaskContext.javaOverride public void rewriteOutputBuffer( // outputBuffer这就是那个“小本本”内存缓存,changeRecord这是监控里看到的一条变化 MapStruct, SourceRecord outputBuffer, SourceRecord changeRecord) { Struct key (Struct) changeRecord.key(); // 获取这条变化的主键比如学号和内容 Struct value (Struct) changeRecord.value(); if (value ! null) { // 看看发生了什么事是新增(CREATE)、修改(UPDATE) 还是 删除(DELETE) Envelope.Operation operation Envelope.Operation.forCode(value.getString(Envelope.FieldName.OPERATION)); switch (operation) { case CREATE: // 【重点】不管是新增还是修改都把最新的样子记下来 // 哪怕是 UPDATE我们也要把它伪装成 READ快照读取因为对下游来说这就是初始数据 case UPDATE: Envelope envelope Envelope.fromSchema(changeRecord.valueSchema()); Struct source value.getStruct(Envelope.FieldName.SOURCE); Struct after value.getStruct(Envelope.FieldName.AFTER);// 取最新的数据 Instant fetchTs Instant.ofEpochMilli((Long) source.get(Envelope.FieldName.TIMESTAMP)); // 重新包装成一条标准的“读取”记录 SourceRecord record new SourceRecord( changeRecord.sourcePartition(), changeRecord.sourceOffset(), changeRecord.topic(), changeRecord.kafkaPartition(), changeRecord.keySchema(), changeRecord.key(), changeRecord.valueSchema(), envelope.read(after, source, fetchTs)); // 【修正操作】 // put 的意思是如果小本本上没有新增就加上如果有修改就覆盖掉旧的 outputBuffer.put(key, record); break; case DELETE: // 【修正操作】 // 如果监控显示这个人走了直接从小本本上划掉他的名字 // 这样下游就根本不知道这个人曾经存在过数据就干净了 outputBuffer.remove(key); break; case READ: // 回填阶段不应该出现 READ 事件如果有就是出 Bug 了 throw new IllegalStateException( String.format( Data change record shouldnt use READ operation, the the record is %s., changeRecord)); } } }第三步发送最终结果等监控录像看完了回填结束内存里的 outputBuffer 就是最准确的数据了。这时候 SeaTunnel 才会把它一次性发给下游。caseMySQL snapshot read - exactly-once快照阶段读不丢不重内存表缓存整个Split数据读取high - low之间的数据按照顺序重放数据到内存表合并key相同的数据保证split high-watermark之前的数据一定被读取到不丢失2增量阶段读 binlog 不丢不重步骤增量阶段读 binlog 不丢不重从 min(split-high-watermark) 开始读取 binlog到 max(split-high-watermark) 结束执行 Exactly-Once 过滤丢弃每个 splitlow~high watermark之间的数据快照阶段已处理保留 split(high-watermark) 之后的数据3.6.3 动态加减表先执行Savepoint停止当前任务修改配置添加或移除目标表恢复任务Sink 端自动识别并增减对应资源3.6.4 CDC 写入 Exactly-Once开启数据库Upsert模式Insert / Update 操作统一转为 Upsert 写入幂等按 key 收集事件保持同一 key 的事件顺序失败/暂停恢复后数据不重复3.6.5 资源优化空闲子任务自动关闭CDC 增量阶段快照读取完成后多余 Task Group 自动关闭仅保留一个 Task Group 负责 binlog 读取。空闲子任务在 binlog 开始读取之后才会关闭保证数据完整性。最终仅保留 Task-Group 1 读写 binlog其余子任务全部关闭大幅节省计算资源。34 任务配置指南4.1 env 配置参数默认值说明parallelism1Source Reader - Sink Writer 的组合数1 Reader 1 Writer 为一组。需 Source 支持分片才有效否则多个并行仅有一个工作。checkpoint.interval300000msCheckpoint 执行间隔。流任务默认 300s不建议配置过小影响性能且易导致失败。Batch 任务默认关闭。出现CheckpointCoordinator inside have error说明保存点超时一般是 reader/writer/committer 处理逻辑卡住需在 seatunnel.yml 中增大 timeout 配置。4.2 Source 分片配置fetch_size、split.size、batch_size需根据文档确认 Source 是否支持分片以及通过哪些配置开启/指定分片方式影响分片的配置分片键 / 分片大小 / 分片个数影响分片效果分片键数据类型 / 数据分布 / 离散稀疏程度即使有分片也不能保证数据均匀读取可能产生数据倾斜配置参数解释调优参考split.size任务怎么切影响并行读取行数split.size 只在使用 table_path 时生效使用query时不生效。小表使用默认值 8096百万级数据可以尝试 10000 ~ 20000千万级数据可以尝试 20000 ~ 50000更大数据量结合并行度、源库压力、单行大小继续压fetch_size数据怎么拉影响单个 Reader 的拉取批次batch_size数据怎么写影响 Sink 的批量写入攒够 batch_size会写一次到了 checkpoint.interval也会写一次。太小写入次数多吞吐可能上不去太大Sink 端缓存更多单次提交压力更大Source 端读取慢优先看parallelism、split.size、fetch_size是否使用 table_pathSink 端写入慢优先看batch_size目标库写入能力是否有主键冲突或upsert目标表索引是否过多。举个例子parallelism 4任务并行度是 4split.size 10000源表按约 10000 行一个 split 切分fetch_size 2000每个 Reader 每次从数据库拉取一批数据batch_size 2000Sink 端攒够 2000 条后批量写入 env { parallelism 4 job.mode BATCH } source { Jdbc { url jdbc:mysql://127.0.0.1:3306/source_db driver com.mysql.cj.jdbc.Driver user root password 123456 table_path source_db.user_source fetch_size 2000 split.size 10000 } } sink { Jdbc { url jdbc:mysql://127.0.0.1:3306/target_db driver com.mysql.cj.jdbc.Driver user root password 123456 table user_target batch_size 2000 } }4.3 Sink 多表配置需根据文档确认 Sink 是否支持多表参考seatunnel-connector-v2-e2e模块下的用例示例支持多表的 Sink 会自动获取 Source 传递的 database/schema/table 名称可在 Sink 路径配置中使用表达式提取源表路径组合拼接前缀后缀${database_name}.${table_name} ${database_name}_xyz.${table_name}_abc ${table_name}_${table_name}_test5 AI ETL助手5.1 目标ST MCP 服务器作为 LLM 和 ST REST API 之间的中间层允许用户使用自然语言提交、监控和管理数据集成作业检索系统监控信息和作业统计信息配置与 SeaTunnel 实例的连接执行复杂的 SeaTunnel 操作无需编写代码或了解底层 API5.2 ST MCP 系统架构交互流程架构图用户使用自然语言与 Claude 等 LLM 交互LLM 使用模型上下文协议与 SeaTunnel MCP 服务器通信MCP 服务器使用专门工具将请求转化为 API 调用SeaTunnelClient 处理与 SeaTunnel REST API 的 HTTP 通信SeaTunnel REST API 与 SeaTunnel 引擎交互以执行操作结果反向流经同一链条最终以自然语言响应的形式返回给用户5.3 ST MCP 核心组件组件说明组件图FastMCP服务器实现模型上下文协议的核心服务器作为 LLM 的通信端点。SeaTunnel客户端SeaTunnel REST API 的封装器用于处理 HTTP 通信、身份验证和数据格式化。MCP工具按功能分类的工具集合用于将客户端方法并使其可被 MCP 服务器访问。CLI界面用于启动和管理 MCP 服务器的命令行界面。

关于本文作者

来自尧图内容编辑团队

尧图内容编辑团队 内容团队

尧图内容编辑团队

本文由尧图网络内容编辑团队执笔。团队由资深项目经理、前端工程师与设计师组成,所有内容均来自亲手交付的真实项目,先讲清问题、再给出可落地的解法。尧图深耕北京网站建设十年,服务过京华建材集团、智造科技等各行业客户,把一线经验沉淀为可复用的行业观察。

  • 十年建站经验,覆盖建材、制造、服务、文创等
  • 项目经理把关选题与事实准确性
  • 工程师与设计师联合撰写专业细节
  • 统一编辑规范,保证文风与排版一致
  • 每月复盘转化数据,迭代选题方向

延伸阅读

相关资讯与近期热门内容

深度阅读推荐

建站决策前值得细读的三篇

网站改版的5个关键决策
2024-08-12

网站改版的5个关键决策

什么时候该改版、改到什么程度、如何避免流量掉光,京华建材集团改版复盘给出答案。

获取专属建站方案

看完文章,把您的行业与预算告诉我们,免费获取一份量身定制的官网建设方案与报价。

立即免费咨询