apache / apache/seatunnel

[Bug] zh:在使用 SeaTunnel CDC(MySQL → Doris)进行数据同步时,任务运行若干小时后突然崩,en:When using SeaTunnel CDC (MySQL → Doris) for data synchronization, the job crashes suddenly after running for several hours.

Open
#9,626 7 comments 0 reactions 0 assignees View on GitHub
bug
Dominant language
Java
Stars
9.7k
Forks
2.4k
Avg merge
3d 9h
Merged PRs (30d)
204

Description

### Search before asking

- [ ] I had searched in the [issues](https://github.com/apache/seatunnel/issues?q=is%3Aissue+label%3A%22bug%22) and found no similar issues.

### What happened

### 🐛 Bug Report
### 描述

在使用 SeaTunnel CDC(MySQL → Doris)进行数据同步时,使用zeta引擎混合模式部署,同步12张表,单表近两千万的的数据,任务运行若干小时后突然崩溃,提示 EventDataDeserializationException,源自 shyiko/mysql-binlog-connector-java。重启任务可以恢复,但过一段时间仍会再次失败。

### 日志部分如下所示:
Caused by: com.github.shyiko.mysql.binlog.event.deserialization.EventDataDeserializationException: Failed to deserialize data of EventHeaderV4{timestamp=1753472097000, eventType=EXT_UPDATE_ROWS, serverId=1476473036, headerLength=19, dataLength=827, nextPosition=41373241, flags=0}]

### 环境信息
SeaTunnel 版本: 2.3.11

### 部署模式: 集群(Cluster)
Source: MySQL CDC (Debezium)
Sink: Doris
使用的 CDC connector: MySQL-CDC
操作系统: CentOS 7 16H 64GB
JDK: Java 11+
Mysql: 腾讯云 mysql8.0

### 问题重现
该任务运行稳定性依赖 binlog 解析,但某些 binlog event 导致 CDC fetcher 崩溃,并整体中断。重启可临时恢复,但仍在几个小时后再次出现。

### 期望行为
希望:

任务在解析 binlog 出现反序列化异常时能跳过该事件,或降级为错误日志记录,不影响整体运行。

支持配置 "skip.broken.events" = true 或其他容错机制。

SeaTunnel 可以通过内部重试机制对这种 binlog 反序列化错误进行隔离。

如社区需要复现数据或进一步分析,可提供脱敏数据样本或运行环境。期待修复建议 🙏

### SeaTunnel Version

SeaTunnel version : 2.3.11

### SeaTunnel Config

```conf
{
source {
MySQL-CDC {
base-url = "jdbc:mysql://:/"
username = "seatunnel"
password = "***"
table-names = [".table1"]
server_id = 5404
startup_mode = "initial"
scan.incremental.snapshot.enabled = true
incremental.snapshot.chunk.size = 3000
chunk-key-column = "s_id"
debezium.skipped.operations = "none"
include.schema.changes = true
debezium.properties = {
"poll.interval.ms" = "600000"
"heartbeat.interval.ms" = "180000"
"max.batch.size" = "4096"
"max.queue.size" = "16384"
}
}
}

sink {
doris {
fenodes = ":8030"
username = "root"
password = "***"
database = ""
table = "table1"
sink.model = "UNIQUE_KEYS"
sink.primary_key = "s_id"
sink.enable-2pc = true
sink.enable-delete = true
batch_size = 2000
interval_ms = 30000
max_retries = 10
doris.config = {
format = "json"
read_json_by_line = "true"
strip_outer_array = "false"
max_filter_ratio = "0.5"
strict_mode = "false"
}
}
}

metrics {
prometheus {
host = "0.0.0.0" # Prometheus指标监听地址,所有网卡
port = 9091 # 监听端口
metrics-name-include = "source_record_count,sink_write_duration" # 指标过滤,只包含这几个
}
}

error_output {
file {
path = "/data/seatunnel_error/table1" # 错误日志文件路径
format = "json" # 错误文件格式
rolling_policy {
max_part_size = "128MB" # 文件最大分片大小
rollover_interval = "1d" # 文件滚动时间间隔,1天
}
}

}
```

### Running Command

```shell
nohup ./bin/seatunnel.sh --config config/sync_sync_table1.conf -e cluster -n sync_sync_table1 > /data/logs/sync_table1.log 2>&1 &

nohup ./bin/seatunnel.sh --config config/sync_sync_table2.conf -e cluster -n sync_sync_table2 > /data/logs/sync_table2.log 2>&1 &

....

nohup ./bin/seatunnel.sh --config config/sync_sync_table12.conf -e cluster -n sync_sync_table11 > /data/logs/sync_table12.log 2>&1 &
```

### Error Exception

```log
2025-07-26 03:35:51,063 INFO [o.a.s.e.c.j.JobMetricsRunner ] [job-metrics-runner-1001064698008829953] -
***********************************************
Job Progress Information
***********************************************
Job Id : 1001064698008829953
Read Count So Far : 26830015
Write Count So Far : 26830015
Average Read Count : 0/s
Average Write Count : 0/s
Last Statistic Time : 2025-07-26 03:34:51
Current Statistic Time : 2025-07-26 03:35:51
***********************************************

2025-07-26 03:36:13,874 INFO [o.a.s.e.c.j.ClientJobProxy ] [main] - Job (1001064698008829953) end with state FAILED
2025-07-26 03:36:13,875 INFO [c.h.c.LifecycleService ] [main] - hz.client_1 [seatunnel] [5.1] HazelcastClient 5.1 (20220228 - 21f20e7) is SHUTTING_DOWN
2025-07-26 03:36:13,894 INFO [.c.i.c.ClientConnectionManager] [main] - hz.client_1 [seatunnel] [5.1] Removed connection to endpoint: [localhost]:5801:55cb8b9f-7481-4205-b2dd-6850b1bd35f8, connection: ClientConnection{alive=false, connectionId=1, channel=NioChannel{/127.0.0.1:54334->localhost/127.0.0.1:5801}, remoteAddress=[localhost]:5801, lastReadTime=2025-07-26 03:36:13.847, lastWriteTime=2025-07-26 03:36:12.725, closedTime=2025-07-26 03:36:13.890, connected server version=5.1}
2025-07-26 03:36:13,894 INFO [c.h.c.LifecycleService ] [main] - hz.client_1 [seatunnel] [5.1] HazelcastClient 5.1 (20220228 - 21f20e7) is CLIENT_DISCONNECTED
2025-07-26 03:36:13,901 INFO [c.h.c.LifecycleService ] [main] - hz.client_1 [seatunnel] [5.1] HazelcastClient 5.1 (20220228 - 21f20e7) is SHUTDOWN
2025-07-26 03:36:13,901 INFO [s.c.s.s.c.ClientExecuteCommand] [main] - Closed SeaTunnel client......
2025-07-26 03:36:13,901 INFO [s.c.s.s.c.ClientExecuteCommand] [main] - Closed metrics executor service ......
2025-07-26 03:36:13,901 ERROR [o.a.s.c.s.SeaTunnel ] [main] -

===============================================================================

2025-07-26 03:36:13,902 ERROR [o.a.s.c.s.SeaTunnel ] [main] - Fatal Error,

2025-07-26 03:36:13,902 ERROR [o.a.s.c.s.SeaTunnel ] [main] - Please submit bug report in https://github.com/apache/seatunnel/issues

2025-07-26 03:36:13,902 ERROR [o.a.s.c.s.SeaTunnel ] [main] - Reason:SeaTunnel job executed failed

2025-07-26 03:36:13,903 ERROR [o.a.s.c.s.SeaTunnel ] [main] - Exception StackTrace:org.apache.seatunnel.core.starter.exception.CommandExecuteException: SeaTunnel job executed failed
at org.apache.seatunnel.core.starter.seatunnel.command.ClientExecuteCommand.execute(ClientExecuteCommand.java:228)
at org.apache.seatunnel.core.starter.SeaTunnel.run(SeaTunnel.java:40)
at org.apache.seatunnel.core.starter.seatunnel.SeaTunnelClient.main(SeaTunnelClient.java:40)
Caused by: org.apache.seatunnel.engine.common.exception.SeaTunnelEngineException: java.lang.RuntimeException: One or more fetchers have encountered exception
at org.apache.seatunnel.connectors.seatunnel.common.source.reader.fetcher.SplitFetcherManager.checkErrors(SplitFetcherManager.java:147)
at org.apache.seatunnel.connectors.seatunnel.common.source.reader.SourceReaderBase.getNextFetch(SourceReaderBase.java:167)
at org.apache.seatunnel.connectors.seatunnel.common.source.reader.SourceReaderBase.pollNext(SourceReaderBase.java:93)
at org.apache.seatunnel.connectors.cdc.base.source.reader.IncrementalSourceReader.pollNext(IncrementalSourceReader.java:119)
at org.apache.seatunnel.engine.server.task.flow.SourceFlowLifeCycle.collect(SourceFlowLifeCycle.java:159)
at org.apache.seatunnel.engine.server.task.SourceSeaTunnelTask.collect(SourceSeaTunnelTask.java:127)
at org.apache.seatunnel.engine.server.task.SeaTunnelTask.stateProcess(SeaTunnelTask.java:165)
at org.apache.seatunnel.engine.server.task.SourceSeaTunnelTask.call(SourceSeaTunnelTask.java:132)
at org.apache.seatunnel.engine.server.TaskExecutionService$BlockingWorker.run(TaskExecutionService.java:694)
at org.apache.seatunnel.engine.server.TaskExecutionService$NamedTaskWrapper.run(TaskExecutionService.java:1023)
at org.apache.seatunnel.api.tracing.MDCRunnable.run(MDCRunnable.java:43)
at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515)
at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)
at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)
at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)
at java.base/java.lang.Thread.run(Thread.java:829)
Caused by: java.lang.RuntimeException: SplitFetcher thread 0 received unexpected exception while polling the records
at org.apache.seatunnel.connectors.seatunnel.common.source.reader.fetcher.SplitFetcher.runOnce(SplitFetcher.java:165)
at org.apache.seatunnel.connectors.seatunnel.common.source.reader.fetcher.SplitFetcher.run(SplitFetcher.java:81)
... 5 more
Caused by: org.apache.kafka.connect.errors.ConnectException: An exception occurred in the change event producer. This connector will be stopped.
at io.debezium.pipeline.ErrorHandler.setProducerThrowable(ErrorHandler.java:50)
at io.debezium.connector.mysql.MySqlStreamingChangeEventSource$ReaderThreadLifecycleListener.onCommunicationFailure(MySqlStreamingChangeEventSource.java:1559)
at com.github.shyiko.mysql.binlog.BinaryLogClient.listenForEventPackets(BinaryLogClient.java:1079)
at com.github.shyiko.mysql.binlog.BinaryLogClient.connect(BinaryLogClient.java:631)
at com.github.shyiko.mysql.binlog.BinaryLogClient$7.run(BinaryLogClient.java:932)
... 1 more
Caused by: io.debezium.DebeziumException: Failed to deserialize data of EventHeaderV4{timestamp=1753472097000, eventType=EXT_UPDATE_ROWS, serverId=1476473036, headerLength=19, dataLength=827, nextPosition=41373241, flags=0}
at io.debezium.connector.mysql.MySqlStreamingChangeEventSource.wrap(MySqlStreamingChangeEventSource.java:1503)
... 5 more
Caused by: com.github.shyiko.mysql.binlog.event.deserialization.EventDataDeserializationException: Failed to deserialize data of EventHeaderV4{timestamp=1753472097000, eventType=EXT_UPDATE_ROWS, serverId=1476473036, headerLength=19, dataLength=827, nextPosition=41373241, flags=0}
at com.github.shyiko.mysql.binlog.event.deserialization.EventDeserializer.deserializeEventData(EventDeserializer.java:341)
at com.github.shyiko.mysql.binlog.event.deserialization.EventDeserializer.nextEvent(EventDeserializer.java:244)
at io.debezium.connector.mysql.MySqlStreamingChangeEventSource$1.nextEvent(MySqlStreamingChangeEventSource.java:273)
at com.github.shyiko.mysql.binlog.BinaryLogClient.listenForEventPackets(BinaryLogClient.java:1051)
... 3 more
Caused by: java.io.EOFException: Failed to read remaining 21 of 177 bytes from position 1681833. Block length: 545. Initial block length: 823.
at com.github.shyiko.mysql.binlog.io.ByteArrayInputStream.fill(ByteArrayInputStream.java:115)
at com.github.shyiko.mysql.binlog.io.ByteArrayInputStream.read(ByteArrayInputStream.java:105)
at com.github.shyiko.mysql.binlog.event.deserialization.AbstractRowsEventDataDeserializer.deserializeBlob(AbstractRowsEventDataDeserializer.java:403)
at com.github.shyiko.mysql.binlog.event.deserialization.AbstractRowsEventDataDeserializer.deserializeCell(AbstractRowsEventDataDeserializer.java:191)
at com.github.shyiko.mysql.binlog.event.deserialization.AbstractRowsEventDataDeserializer.deserializeRow(AbstractRowsEventDataDeserializer.java:143)
at com.github.shyiko.mysql.binlog.event.deserialization.UpdateRowsEventDataDeserializer.deserializeRows(UpdateRowsEventDataDeserializer.java:71)
at com.github.shyiko.mysql.binlog.event.deserialization.UpdateRowsEventDataDeserializer.deserialize(UpdateRowsEventDataDeserializer.java:58)
at com.github.shyiko.mysql.binlog.event.deserialization.UpdateRowsEventDataDeserializer.deserialize(UpdateRowsEventDataDeserializer.java:33)
at com.github.shyiko.mysql.binlog.event.deserialization.EventDeserializer.deserializeEventData(EventDeserializer.java:335)
... 6 more

at org.apache.seatunnel.core.starter.seatunnel.command.ClientExecuteCommand.execute(ClientExecuteCommand.java:220)
... 2 more

2025-07-26 03:36:13,903 ERROR [o.a.s.c.s.SeaTunnel ] [main] -
===============================================================================

Exception in thread "main" org.apache.seatunnel.core.starter.exception.CommandExecuteException: SeaTunnel job executed failed
at org.apache.seatunnel.core.starter.seatunnel.command.ClientExecuteCommand.execute(ClientExecuteCommand.java:228)
at org.apache.seatunnel.core.starter.SeaTunnel.run(SeaTunnel.java:40)
at org.apache.seatunnel.core.starter.seatunnel.SeaTunnelClient.main(SeaTunnelClient.java:40)
Caused by: org.apache.seatunnel.engine.common.exception.SeaTunnelEngineException: java.lang.RuntimeException: One or more fetchers have encountered exception
at org.apache.seatunnel.connectors.seatunnel.common.source.reader.fetcher.SplitFetcherManager.checkErrors(SplitFetcherManager.java:147)
at org.apache.seatunnel.connectors.seatunnel.common.source.reader.SourceReaderBase.getNextFetch(SourceReaderBase.java:167)
at org.apache.seatunnel.connectors.seatunnel.common.source.reader.SourceReaderBase.pollNext(SourceReaderBase.java:93)
at org.apache.seatunnel.connectors.cdc.base.source.reader.IncrementalSourceReader.pollNext(IncrementalSourceReader.java:119)
at org.apache.seatunnel.engine.server.task.flow.SourceFlowLifeCycle.collect(SourceFlowLifeCycle.java:159)
at org.apache.seatunnel.engine.server.task.SourceSeaTunnelTask.collect(SourceSeaTunnelTask.java:127)
at org.apache.seatunnel.engine.server.task.SeaTunnelTask.stateProcess(SeaTunnelTask.java:165)
at org.apache.seatunnel.engine.server.task.SourceSeaTunnelTask.call(SourceSeaTunnelTask.java:132)
at org.apache.seatunnel.engine.server.TaskExecutionService$BlockingWorker.run(TaskExecutionService.java:694)
at org.apache.seatunnel.engine.server.TaskExecutionService$NamedTaskWrapper.run(TaskExecutionService.java:1023)
at org.apache.seatunnel.api.tracing.MDCRunnable.run(MDCRunnable.java:43)
at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515)
at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)
at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)
at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)
at java.base/java.lang.Thread.run(Thread.java:829)
Caused by: java.lang.RuntimeException: SplitFetcher thread 0 received unexpected exception while polling the records
at org.apache.seatunnel.connectors.seatunnel.common.source.reader.fetcher.SplitFetcher.runOnce(SplitFetcher.java:165)
at org.apache.seatunnel.connectors.seatunnel.common.source.reader.fetcher.SplitFetcher.run(SplitFetcher.java:81)
... 5 more
Caused by: org.apache.kafka.connect.errors.ConnectException: An exception occurred in the change event producer. This connector will be stopped.
at io.debezium.pipeline.ErrorHandler.setProducerThrowable(ErrorHandler.java:50)
at io.debezium.connector.mysql.MySqlStreamingChangeEventSource$ReaderThreadLifecycleListener.onCommunicationFailure(MySqlStreamingChangeEventSource.java:1559)
at com.github.shyiko.mysql.binlog.BinaryLogClient.listenForEventPackets(BinaryLogClient.java:1079)
at com.github.shyiko.mysql.binlog.BinaryLogClient.connect(BinaryLogClient.java:631)
at com.github.shyiko.mysql.binlog.BinaryLogClient$7.run(BinaryLogClient.java:932)
... 1 more
Caused by: io.debezium.DebeziumException: Failed to deserialize data of EventHeaderV4{timestamp=1753472097000, eventType=EXT_UPDATE_ROWS, serverId=1476473036, headerLength=19, dataLength=827, nextPosition=41373241, flags=0}
at io.debezium.connector.mysql.MySqlStreamingChangeEventSource.wrap(MySqlStreamingChangeEventSource.java:1503)
... 5 more
Caused by: com.github.shyiko.mysql.binlog.event.deserialization.EventDataDeserializationException: Failed to deserialize data of EventHeaderV4{timestamp=1753472097000, eventType=EXT_UPDATE_ROWS, serverId=1476473036, headerLength=19, dataLength=827, nextPosition=41373241, flags=0}
at com.github.shyiko.mysql.binlog.event.deserialization.EventDeserializer.deserializeEventData(EventDeserializer.java:341)
at com.github.shyiko.mysql.binlog.event.deserialization.EventDeserializer.nextEvent(EventDeserializer.java:244)
at io.debezium.connector.mysql.MySqlStreamingChangeEventSource$1.nextEvent(MySqlStreamingChangeEventSource.java:273)
at com.github.shyiko.mysql.binlog.BinaryLogClient.listenForEventPackets(BinaryLogClient.java:1051)
... 3 more
Caused by: java.io.EOFException: Failed to read remaining 21 of 177 bytes from position 1681833. Block length: 545. Initial block length: 823.
at com.github.shyiko.mysql.binlog.io.ByteArrayInputStream.fill(ByteArrayInputStream.java:115)
at com.github.shyiko.mysql.binlog.io.ByteArrayInputStream.read(ByteArrayInputStream.java:105)
at com.github.shyiko.mysql.binlog.event.deserialization.AbstractRowsEventDataDeserializer.deserializeBlob(AbstractRowsEventDataDeserializer.java:403)
at com.github.shyiko.mysql.binlog.event.deserialization.AbstractRowsEventDataDeserializer.deserializeCell(AbstractRowsEventDataDeserializer.java:191)
at com.github.shyiko.mysql.binlog.event.deserialization.AbstractRowsEventDataDeserializer.deserializeRow(AbstractRowsEventDataDeserializer.java:143)
at com.github.shyiko.mysql.binlog.event.deserialization.UpdateRowsEventDataDeserializer.deserializeRows(UpdateRowsEventDataDeserializer.java:71)
at com.github.shyiko.mysql.binlog.event.deserialization.UpdateRowsEventDataDeserializer.deserialize(UpdateRowsEventDataDeserializer.java:58)
at com.github.shyiko.mysql.binlog.event.deserialization.UpdateRowsEventDataDeserializer.deserialize(UpdateRowsEventDataDeserializer.java:33)
at com.github.shyiko.mysql.binlog.event.deserialization.EventDeserializer.deserializeEventData(EventDeserializer.java:335)
... 6 more

at org.apache.seatunnel.core.starter.seatunnel.command.ClientExecuteCommand.execute(ClientExecuteCommand.java:220)
... 2 more
2025-07-26 03:36:13,914 INFO [s.c.s.s.c.ClientExecuteCommand] [SeaTunnel-CompletableFuture-Thread-0] - run shutdown hook because get close signal
```

### Zeta or Flink or Spark Version

相关连接器版本
connector-doris-2.3.11.jar
connector-cdc-mysql-2.3.11.jar

### Java or Scala Version

jdk11

### Screenshots
## seatunnel web ui
Image

## grafana
这是grafana+prometheus 监控到的截图,可见目前只有两个任务正在正常运行,其余10个任务都失败了
Image

### Are you willing to submit PR?

- [x] Yes I am willing to submit a PR!

### Code of Conduct

- [x] I agree to follow this project's [Code of Conduct](https://www.apache.org/foundation/policies/conduct)

Contributor guide

No contributing guide indexed for this repository

Research direction

Start with the stack-trace entry points SplitFetcherManager.checkErrors, MySqlStreamingChangeEventSource, BinaryLogClient, and EventDeserializer, then reproduce the failure with the supplied MySQL-CDC configuration and EXT_UPDATE_ROWS log. Done means the malformed binlog event behavior is defined and verified without unexpectedly terminating the CDC job.

Written by the indexing model from the issue text.

Assessment

Tech stack
java, mysql
Domain
data-engineering, databases, distributed-systems
Issue type
Bug
Difficulty
5/5
Estimated time
Over a week
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
30/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.