无法在S3中写入Iceberg记录
我在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导航站,无授权禁止任何主体转载、抄袭、复制内容,亦不得私自架设镜像站点。一经侵权,本站将通过法律途径追责。