Flink CDC+Hudi同步MySQL至S3出重复行,疑与状态TTL清除有关?
问题描述
使用Flink 1.16.2,通过Flink CDC + Apache Hudi将MySQL数据同步至AWS S3,作业代码如下:
parallelism = 1 env = StreamExecutionEnvironment.get_execution_environment(config) env.set_parallelism(parallelism) # I don't know what's that, try 2 env.enable_checkpointing(10 * 60 * 1000) # milliseconds # checkpoints have to complete within one minute, or are discarded env.get_checkpoint_config().set_checkpoint_timeout(int(10 * 60 * 1000)) env.get_checkpoint_config().set_checkpointing_mode( CheckpointingMode.EXACTLY_ONCE ) env.disable_operator_chaining() # If drop this line, create multiple pipeline in one python job will raise error in Flink GUI settings = EnvironmentSettings.new_instance().in_streaming_mode().build() t_env = StreamTableEnvironment.create(env, environment_settings=settings) ss = t_env.create_statement_set() source_sql = f""" CREATE TABLE mysql_source ( db_name STRING METADATA FROM 'database_name' VIRTUAL, table_name STRING METADATA FROM 'table_name' VIRTUAL, id BIGINT, content STRING, update_time TIMESTAMP(3) ) WITH ( 'connector' = 'mysql-cdc', 'hostname' = '{db_secret["host"]}', 'port' = '{db_secret["port"]}', 'username' = '{db_secret["username"]}', 'password' = '{db_secret["password"]}', 'database-name' = 'my_db', 'table-name' = 'my_table', 'server-id' = '{server_id}', 'debezium.snapshot.mode' = 'schema_only', 'scan.startup.mode' = 'timestamp', 'scan.startup.timestamp-millis' = '{binlog_start_time}', 'scan.incremental.snapshot.enabled' = 'true', 'scan.incremental.snapshot.chunk.key-column' = 'id' ); """ t_env.execute_sql(source_sql) sink_sql = f""" CREATE TABLE hudi_sink ( id BIGINT, content STRING, update_time TIMESTAMP(3), hudi_ts double, PRIMARY KEY (id) NOT ENFORCED ) WITH ( 'connector' = 'hudi', 'path' = 's3a://xxxx', 'table.type' = 'COPY_ON_WRITE', 'write.precombine.field' = 'hudi_ts', 'write.operation' = 'upsert', 'hoodie.datasource.write.recordkey.field' = 'id', 'hoodie.datasource.write.partitionpath.field' = '', 'write.tasks' = '1', 'compaction.tasks' = '1', 'clean.retain_commits' = '6', 'hoodie.keep.min.commits' = '7', 'hoodie.keep.max.commits' = '8', 'compaction.async.enabled' = 'true', 'hoodie.write.lock.provider' = 'org.apache.hudi.client.transaction.lock.InProcessLockProvider', 'index.partition.regex' = 'false', 'index.bootstrap.enabled' = 'false', 'write.rate.limit' = '2000' ); """ t_env.execute_sql(sink_sql) insert_sql = f""" INSERT INTO hudi_sink SELECT id,content,update_time, UNIX_TIMESTAMP() as hudi_ts FROM mysql_source where db_name = 'my_db' and table_name = 'my_table'; """ ss.add_insert_sql(insert_sql) ss.attach_as_datastream() env.execute(f'mysql_cdc_to_hudi')
更新MySQL中的某一行后,Hudi表中出现两条同id的记录,主键未生效,无法更新现有记录而是插入新行。同时发现如下日志:
2025-09-01 01:10:20,612 INFO org.apache.flink.table.runtime.operators.sink.SinkUpsertMaterializer [] - The state is cleared because of state ttl. This will result in incorrect result. You can increase the state ttl to avoid this.
请问该日志提示的状态TTL清除是否与我的问题相关?
解答
这条日志提示的状态TTL清除确实和你的Hudi主键不生效、重复插入的问题直接相关,原因如下:
- SinkUpsertMaterializer是Flink处理Upsert逻辑的核心算子,它会维护一个状态来记录主键对应的最新操作类型(插入/更新)。当状态TTL到期被清除后,算子会丢失之前存储的主键状态信息。
- 当MySQL的更新事件到来时,算子无法识别这是一条现有记录的更新,会错误地将其当作新插入操作发送给Hudi;而Hudi的Upsert逻辑依赖主键状态判断最新版本,最终导致同ID的重复记录出现。
解决步骤
- 调整状态TTL参数:在创建
StreamTableEnvironment后添加状态TTL配置,确保时长大于你的最大数据延迟时间(比如设置为1天):
t_env.get_config().set("table.exec.state.ttl", "86400000") # 单位:毫秒,对应1天
- 匹配Checkpoint与TTL时长:你的Checkpoint间隔是10分钟,需保证状态TTL时长大于Checkpoint间隔,避免状态还未被持久化就被清除。
- 优化Hudi预合并字段:当前用
UNIX_TIMESTAMP()生成hudi_ts是Flink处理数据的时间,建议改用MySQL源表的update_time生成该字段,让Hudi更准确判断最新版本:
INSERT INTO hudi_sink SELECT id,content,update_time, UNIX_TIMESTAMP(update_time) as hudi_ts FROM mysql_source where db_name = 'my_db' and table_name = 'my_table';
也可以直接将Hudi sink的write.precombine.field改为update_time(需保证字段类型匹配)。
额外排查点
- 确认Hudi配置的
hoodie.datasource.write.recordkey.field确实为id,与表定义的主键一致; - 检查Flink作业并行度与Hudi的
write.tasks是否匹配,避免并行度不一致导致的状态混乱; - 验证MySQL CDC连接器是否正确发送UPDATE类型的事件(可通过打印日志查看事件类型)。
内容的提问来源于stack exchange,提问作者Rinze
相关产品推荐
相关产品推荐

