无法在S3中写入Iceberg记录

后端开发 2026-07-10

我在EMR上使用Flink,Kafka安装在一个EKS集群中,Flink可以与Kafka正常通信。连接没有问题。

我想使用Table API,通过Flink的 Sink表把来自Kafka主题的记录写入到S3。然而,尽管数据被写入,但数据被某种方式损坏。

用于创建表的代码,以及相应触发的命令。

CREATE OR REPLACE TABLE kafka_source (  
  id    INT,  
  data  STRING  
) WITH (  
  'connector'                    = 'kafka',  
  'topic'                        = 'test-topic',  
  'properties.bootstrap.servers' = 'something.us-east-1.elb.amazonaws.com:9094',  
  'properties.group.id'          = 'flink-iceberg-consumer',  
  'scan.startup.mode'            = 'earliest-offset',  
  'format'                       = 'json'  
);  

CREATE CATALOG glue_catalog WITH (  
  'type'         = 'iceberg',  
  'warehouse'    = 's3://eks-benchmark-iceberg/warehouse/',  
  'catalog-impl' = 'org.apache.iceberg.aws.glue.GlueCatalog',  
  'io-impl'      = 'org.apache.iceberg.aws.s3.S3FileIO' 
);  

USE CATALOG glue_catalog;  

CREATE DATABASE IF NOT EXISTS benchmark_db;  

USE benchmark_db;  

CREATE TABLE IF NOT EXISTS orders_iceberg (  
  id    INT,  
  data  STRING  
);  

USE CATALOG default_catalog;
SET 'execution.runtime-mode' = 'streaming';
ADD JAR '/usr/lib/flink/lib/flink-sql-connector-kafka-3.3.0-1.20.jar';  
INSERT INTO glue_catalog.benchmark_db.orders_iceberg 
SELECT id, data FROM kafka_source;

现在,我在Flink UI里观察到两个问题。第一,虽然只摄取了一个记录,但显示处理的记录数超过了950条。此外,记录也被损坏。我也无法查看该记录。我附上了几张Flink UI的截图:

主页详情输出端详情检查点详情

没有任何异常或错误。我已经检查过。我正在使用下面的记录来产生和消费。

{"id": 1, "data": "hello"}

Flink version: 1.20.0
Kafka version: 4.x

那么我到底错在了哪里?

解决方案

竟然是一件挺有意思的事。

跑完爬虫之后,我就能在Athena里取到数据。莫名其妙地,S3 Select不允许查询。

代码或做法本身没有问题。

不过仍需查清楚为何尽管我只通过生产者发送了1 条记录,表中却能看到如此多的记录被发送。

站内所有文章版权归属LeftHeroAI导航站,无授权禁止任何主体转载、抄袭、复制内容,亦不得私自架设镜像站点。一经侵权,本站将通过法律途径追责。

相关文章