PyFlink/Flink Table API写入S3上Apache Iceberg无数据输出问题排查
问题:PyFlink Iceberg实时同步无数据写入
环境信息
- Flink镜像:
flink:1.16.1-scala_2.12 - 依赖Jar包:
- bundle-2.20.18.jar
- flink-shaded-hadoop-2-uber-2.8.3-10.0.jar
- flink-sql-connector-hive-2.3.9_2.12-1.16.1.jar
- guava-30.1-jre.jar
- hadoop-common-2.8.3.jar
- iceberg-flink-runtime-1.16-1.3.1.jar
操作与现象
- 启动
jobmanager.sh和taskmanager.sh - 执行PyFlink脚本,脚本成功创建目标Iceberg表
- 日志显示能扫描源Iceberg表并打印最新
manifest.json名称,但无数据写入目标表 - CPU持续占用约50%,已排除资源不足(M1设备分配全部资源)
疑问
不确定INSERT ... SELECT语法在Iceberg实时同步场景下是否可用,怀疑INSERT语句存在问题。
相关代码
settings = EnvironmentSettings.new_instance().in_streaming_mode().build() t_env = TableEnvironment.create(environment_settings=settings) catalog_name = "glue_catalog" staging_database_name = "source" dest_database_name = "dest" staging_table = "source_iceberg_partitioned_hourly" dest_table = "new_iceberg_table" warehouse_path = "s3:path....." t_env.execute_sql( f""" CREATE CATALOG {catalog_name} WITH ( 'type'='iceberg', 'catalog-impl'='org.apache.iceberg.aws.glue.GlueCatalog', 'warehouse'='{warehouse_path}', 'aws.region'='eu-central-1', 'io-impl'='org.apache.iceberg.aws.s3.S3FileIO' ) """ ) t_env.execute_sql(f"USE CATALOG {catalog_name};") t_env.execute_sql(f"CREATE DATABASE IF NOT EXISTS {staging_database_name}") t_env.execute_sql(f"CREATE DATABASE IF NOT EXISTS {dest_database_name}") t_env.execute_sql(f"USE {dest_database_name};") staging_table_full = f"{catalog_name}.{staging_database_name}.{staging_table}" dest_table_full = f"{catalog_name}.{dest_database_name}.{dest_table}" create_table_statement = f""" CREATE TABLE IF NOT EXISTS {dest_table} ( -- 列定义省略 ) PARTITIONED BY (collector_minute) WITH ( 'type'='iceberg', 'format-version'='2', 'write.format.default'='parquet', 'write.metadata.delete-after-commit.enabled'='true', 'write.distribution-mode'='hash', 'path'='{warehouse_path}' ) """ t_env.execute_sql(create_table_statement) latest_snapshot_id = "12345678901234567890" t_env.execute_sql( f""" INSERT INTO {dest_table_full} SELECT *, MINUTE(collector_tstamp) AS collector_minute FROM {staging_table_full} /*+ OPTIONS('streaming' = 'true', 'monitor-interval' = '10s', 'start-snapshot-id' = '{latest_snapshot_id}') */ """ )
问题分析与解决建议
语法变量错误:
你的代码中INSERT语句最初写的是{query_table_full},但实际定义的目标表变量是dest_table_full,这会导致写入目标不存在或错误,修正为{dest_table_full}即可。Iceberg流读参数验证:
- 确认
start-snapshot-id对应的快照存在且包含数据,如果该快照为空,流任务只会监听后续新增快照,不会读取历史数据。可临时移除该参数,测试是否能读取全量数据再监听增量。 monitor-interval设为10s合法,但可尝试调整为30s减少扫描频率,避免不必要的资源消耗。
- 确认
目标表路径配置问题:
目标表的path参数直接设为仓库根路径{warehouse_path},会导致表数据与其他表冲突,应指定具体表路径,比如'{warehouse_path}/{dest_database_name}/{dest_table}'。字段类型匹配检查:
确认目标表中collector_minute的类型与MINUTE(collector_tstamp)的返回值类型一致(比如collector_tstamp是TIMESTAMP类型时,MINUTE()返回INT,目标表该字段需定义为INT)。任务状态排查:
- 访问Flink Web UI(默认8081端口)查看算子链、数据流向,检查是否有数据卡在某个算子环节。
- 查看TaskManager日志,搜索
IcebergSink相关条目,确认是否存在写入失败、数据过滤等细节。
依赖兼容性验证:
确认iceberg-flink-runtime-1.16-1.3.1.jar与Flink 1.16.1版本兼容,同时检查Guava版本是否存在冲突(Iceberg依赖的Guava版本需与你引入的30.1-jre兼容)。
内容的提问来源于stack exchange,提问作者amarius
相关产品推荐
相关产品推荐

