使用Flink Kubernetes Operator通过Flink SQL向Kafka写数据遇阻
问题:Flink SQL消费Kafka数据后无输出到Upsert-Kafka主题
我运行以下Flink SQL代码,从Kafka的demo主题读取数据并插入到upsert-kafka类型的demosink主题中:
CREATE TABLE IF NOT EXISTS some_source_table ( myField1 VARCHAR, myField2 VARCHAR ) WITH ( 'connector' = 'kafka', 'topic' = 'demo', 'properties.bootstrap.servers' = '***', 'properties.group.id' = 'some-id-1', 'scan.startup.mode' = 'latest-offset', 'format' = 'json', 'json.timestamp-format.standard' = 'ISO-8601', 'scan.topic-partition-discovery.interval'= '60000', 'json.fail-on-missing-field' = 'false', 'json.ignore-parse-errors' = 'true', 'properties.security.protocol' = 'SASL_SSL', 'properties.sasl.mechanism' = 'PLAIN', 'properties.sasl.jaas.config' = 'org.apache.kafka.common.security.plain.PlainLoginModule required username=*** password=***;', 'properties.ssl.endpoint.identification.algorithm' = 'https' ); CREATE TABLE IF NOT EXISTS some_sink_table ( myField1 VARCHAR, myField2 VARCHAR, PRIMARY KEY (`myField1`) NOT ENFORCED ) WITH ( 'connector' = 'upsert-kafka', 'topic' = 'demosink', 'properties.bootstrap.servers' = '***', 'key.format' = 'json', 'value.format' = 'json', 'properties.security.protocol' = 'SASL_SSL', 'properties.sasl.mechanism' = 'PLAIN', 'properties.sasl.jaas.config' = 'org.apache.kafka.common.security.plain.PlainLoginModule required username=*** password=***;', 'properties.ssl.endpoint.identification.algorithm' = 'https' ); INSERT INTO some_sink_table SELECT * FROM some_source_table;
在Confluent平台能看到源数据已被消费,但demosink主题始终没有数据产出。我使用的Flink版本为1.19.1,参考相关示例添加了以下依赖:
<dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-kafka</artifactId> <version>3.1.0-1.18</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-sql-connector-kafka</artifactId> <version>3.1.0-1.18</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-scala_2.12</artifactId> <version>${flink.version}</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-clients</artifactId> <version>${flink.version}</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-table-planner_2.12</artifactId> <version>${flink.version}</version> <scope>provided</scope> </dependency>
我也尝试过将flink-sql-connector-kafka和flink-clients的jar包复制到Dockerfile的lib目录中,想请教是否存在依赖缺失或不兼容的问题?
问题分析与解决方案
1. 核心问题:依赖版本不兼容
你使用的Flink版本是1.19.1,但Kafka连接器(flink-connector-kafka、flink-sql-connector-kafka)用的是3.1.0-1.18,该版本是适配Flink 1.18的,与1.19.1存在版本不兼容。Flink连接器版本必须与核心版本严格匹配,版本错位会导致序列化逻辑异常、连接器内部流程中断,进而引发数据无法正常写入sink主题。
2. 依赖冗余问题
flink-sql-connector-kafka已经包含了flink-connector-kafka的全部功能,同时引入两者会造成依赖冲突,干扰类加载逻辑。
3. 修复步骤
- 统一依赖版本:将Kafka连接器版本替换为适配Flink 1.19.1的
3.2.0-1.19(Flink 1.19对应的Kafka连接器版本为3.2.0)。 - 移除冗余依赖:删除
flink-connector-kafka依赖,仅保留flink-sql-connector-kafka。 - 检查依赖范围:
flink-table-planner_2.12的scope设为provided,需确认Flink Kubernetes Operator集群环境已包含该依赖;若集群未提供,需将scope改为compile,或在Docker镜像中添加该jar包。 - 清理Docker镜像:确保镜像中仅保留正确版本的连接器jar包,删除旧版本残留文件。
额外排查方向
- 主键有效性:确认源数据中
myField1(sink表主键)不为空,Upsert-Kafka会基于主键做更新/插入,主键为空会导致数据被丢弃。 - 任务日志排查:查看TaskManager日志,检查是否存在序列化错误、Kafka写入权限不足、数据格式解析失败等异常信息。
- 主题权限验证:确认Flink任务使用的Kafka账号拥有
demosink主题的写入权限。
内容的提问来源于stack exchange,提问作者user28527275
相关产品推荐
相关产品推荐

