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

PyFlink/Flink Table API写入S3上Apache 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

操作与现象

  1. 启动jobmanager.sh和taskmanager.sh
  2. 执行PyFlink脚本,脚本成功创建目标Iceberg表
  3. 日志显示能扫描源Iceberg表并打印最新manifest.json名称,但无数据写入目标表
  4. 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}') */
    """
)

问题分析与解决建议

  1. 语法变量错误:
    你的代码中INSERT语句最初写的是{query_table_full},但实际定义的目标表变量是dest_table_full,这会导致写入目标不存在或错误,修正为{dest_table_full}即可。

  2. Iceberg流读参数验证:

    • 确认start-snapshot-id对应的快照存在且包含数据,如果该快照为空,流任务只会监听后续新增快照,不会读取历史数据。可临时移除该参数,测试是否能读取全量数据再监听增量。
    • monitor-interval设为10s合法,但可尝试调整为30s减少扫描频率,避免不必要的资源消耗。
  3. 目标表路径配置问题:
    目标表的path参数直接设为仓库根路径{warehouse_path},会导致表数据与其他表冲突,应指定具体表路径,比如'{warehouse_path}/{dest_database_name}/{dest_table}'。

  4. 字段类型匹配检查:
    确认目标表中collector_minute的类型与MINUTE(collector_tstamp)的返回值类型一致(比如collector_tstamp是TIMESTAMP类型时,MINUTE()返回INT,目标表该字段需定义为INT)。

  5. 任务状态排查:

    • 访问Flink Web UI(默认8081端口)查看算子链、数据流向,检查是否有数据卡在某个算子环节。
    • 查看TaskManager日志,搜索IcebergSink相关条目,确认是否存在写入失败、数据过滤等细节。
  6. 依赖兼容性验证:
    确认iceberg-flink-runtime-1.16-1.3.1.jar与Flink 1.16.1版本兼容,同时检查Guava版本是否存在冲突(Iceberg依赖的Guava版本需与你引入的30.1-jre兼容)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 09:24:52