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

使用Flink Kubernetes Operator通过Flink SQL向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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 23:54:50