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

Flink迁移MySQL至Iceberg时TaskManager崩溃问题求助

问题

在使用Flink将约4TB的MySQL数据流式迁移至Iceberg时,TaskManager出现崩溃,报错org.apache.flink.runtime.jobmaster.JobMasterException: TaskManager with id localhost:26211-081070 is no longer reachable.,同时伴随Java堆空间不足异常。即使仅执行SELECT查询流式处理数据,处理约1万条记录后也会触发相同错误。怀疑INSERT INTO xx SELECT * FROM xx会一次性加载全表到内存,但单独查询也出错,不清楚原因,也不确定是否需要手动清理内存。

报错日志

23/10/17 14:43:34 ERROR Executor: Exception in task 0.0 in stage 0.0 (TID 0)/ 1]
java.sql.SQLException: Java heap space
        at com.mysql.cj.jdbc.exceptions.SQLError.createSQLException(SQLError.java:130)
        at com.mysql.cj.jdbc.exceptions.SQLExceptionsMapping.translateException(SQLExceptionsMapping.java:122)
        at com.mysql.cj.jdbc.ClientPreparedStatement.executeInternal(ClientPreparedStatement.java:916)
        at com.mysql.cj.jdbc.ClientPreparedStatement.executeQuery(ClientPreparedStatement.java:972)
        at org.apache.spark.sql.execution.datasources.jdbc.JDBCRDD.compute(JDBCRDD.scala:314)
        ...(省略重复日志内容)
ConnectionRefusedError: [Errno 111] Connection refused

代码实现

env = StreamExecutionEnvironment.get_execution_environment()
env.set_parallelism(1).enable_checkpointing(3000)
table_env = StreamTableEnvironment.create(env)

# source table
table_env.execute_sql("""
    CREATE TABLE src_result (
        id STRING,
        input_md5 STRING,
        output_md5 STRING,
        log STRING,
        metric STRING,
        create_time TIMESTAMP(6),      
        update_time TIMESTAMP(6),
        workflow_id_id STRING,
        error_details STRING,
        error_stage STRING,
        error_type STRING
    ) WITH (
        'connector' = 'jdbc',
        'url' = 'jdbc:mysql://xxx',
        'table-name' = 'result',
        'username' = 'ars_dev',
        'password' = '01234567'
    );
""")

# target iceberg
table_env.execute_sql("CREATE CATALOG iceberg WITH ("
                    "'type'='iceberg', "
                    "'catalog-type'='hive', "
                    "'uri'='thrift://xxx',"
                    "'warehouse'='xxx')")

# start migration
t_res = table_env.execute_sql("""
    INSERT INTO iceberg.ars.results
    SELECT 
        id,
        input_md5,
        output_md5,
        log,
        metric,
        CAST(create_time AS STRING),
        CAST(update_time AS STRING),
        workflow_id_id,
        error_details,
        error_stage,
        error_type  
    FROM src_result
""")

分析与解决方案

核心原因

  1. JDBC Connector未配置分页读取:默认情况下,Flink JDBC Source会尝试一次性拉取全量数据到内存,4TB大表直接触发内存溢出,即使单独SELECT也会因无分页机制,累积数据后耗尽内存。
  2. 并行度过低:代码中并行度设为1,所有数据集中在单个TaskManager处理,内存压力集中,极易触发堆内存不足。
  3. 大字段占用内存:表中log、metric等字段可能存储大文本数据,单条记录内存占用高,快速累积后耗尽堆空间。

解决步骤

  1. 配置JDBC分页与分片参数
    在MySQL Source的WITH配置中添加分页和分片参数,让Flink分批、分片读取数据:

    'connector' = 'jdbc',
    'url' = 'jdbc:mysql://xxx',
    'table-name' = 'result',
    'username' = 'ars_dev',
    'password' = '01234567',
    'scan.fetch-size' = '1000', -- 每次从数据库拉取的行数
    'scan.partition.column' = 'id', -- 用于分片的有序列(如主键、自增ID)
    'scan.partition.lower-bound' = '1', -- 分片列最小值
    'scan.partition.upper-bound' = '10000000', -- 分片列最大值
    'scan.partition.num' = '10' -- 分片数量,根据集群资源调整
    

    注意:scan.partition.column必须是有序列,保证数据均匀分片,避免单分片数据量过大。

  2. 调整并行度与TaskManager内存

    • 提高并行度:根据集群TaskManager数量和CPU核心数,设置合理并行度(如8、16),分散处理压力:
      env.set_parallelism(8)
      
    • 增大TaskManager堆内存:在flink-conf.yaml中调整taskmanager.memory.process.size或taskmanager.memory.task.heap.size,比如设置为16g或更高,匹配数据量需求。
  3. 优化大字段处理

    • 若log、metric为大文本,可将其单独存储到对象存储(如OSS、S3),Iceberg仅存储文件路径,减少单条记录内存占用;
    • 迁移时过滤非必要大字段,只同步核心数据。
  4. 无需手动清理内存
    Flink内存由框架自动管理,无需手动清理。内存问题本质是数据读取和处理配置不合理,调整上述配置即可解决。

内容的提问来源于stack exchange,提问作者Martin_5125

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 01:49:55