You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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的重复记录出现。

解决步骤

  1. 调整状态TTL参数:在创建StreamTableEnvironment后添加状态TTL配置,确保时长大于你的最大数据延迟时间(比如设置为1天):
t_env.get_config().set("table.exec.state.ttl", "86400000")  # 单位:毫秒,对应1天
  1. 匹配Checkpoint与TTL时长:你的Checkpoint间隔是10分钟,需保证状态TTL时长大于Checkpoint间隔,避免状态还未被持久化就被清除。
  2. 优化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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.12 10:44:52