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 """)
分析与解决方案
核心原因
- JDBC Connector未配置分页读取:默认情况下,Flink JDBC Source会尝试一次性拉取全量数据到内存,4TB大表直接触发内存溢出,即使单独SELECT也会因无分页机制,累积数据后耗尽内存。
- 并行度过低:代码中并行度设为1,所有数据集中在单个TaskManager处理,内存压力集中,极易触发堆内存不足。
- 大字段占用内存:表中
log、metric等字段可能存储大文本数据,单条记录内存占用高,快速累积后耗尽堆空间。
解决步骤
配置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必须是有序列,保证数据均匀分片,避免单分片数据量过大。调整并行度与TaskManager内存
- 提高并行度:根据集群TaskManager数量和CPU核心数,设置合理并行度(如8、16),分散处理压力:
env.set_parallelism(8) - 增大TaskManager堆内存:在
flink-conf.yaml中调整taskmanager.memory.process.size或taskmanager.memory.task.heap.size,比如设置为16g或更高,匹配数据量需求。
- 提高并行度:根据集群TaskManager数量和CPU核心数,设置合理并行度(如8、16),分散处理压力:
优化大字段处理
- 若
log、metric为大文本,可将其单独存储到对象存储(如OSS、S3),Iceberg仅存储文件路径,减少单条记录内存占用; - 迁移时过滤非必要大字段,只同步核心数据。
- 若
无需手动清理内存
Flink内存由框架自动管理,无需手动清理。内存问题本质是数据读取和处理配置不合理,调整上述配置即可解决。
内容的提问来源于stack exchange,提问作者Martin_5125
相关产品推荐
相关产品推荐

