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

Flink Table API读取某Kafka Topic时列值错乱问题排查求助

我用Scala开发了一个Flink作业,负责关联多个Kafka Topic流做数据增强。通过Table API指定Schema读取这些Topic后转成DataStream,但其中一个Topic转成DataStream后打印日志时出现列值错乱(比如第1列始终显示第4列的值),不过Schema的列名称是正确的。

已排查内容

  • Schema读取正确
  • 列未按字母顺序排列
  • Kafka Topic中的消息内容符合预期
  • 相同代码读取其他Topic均正常
  • 已升级到最新版Flink

临时解决方案

  • 在select语句中手动调整列顺序
  • 直接改用DataStream API读取该流

补充说明与代码

使用Avro Schema,Kafka Topic格式为avro-confluent,构建Table Descriptor的代码如下:

val schemaString   = Source.fromResource(schemaPath).mkString
val schemaDataType = AvroSchemaConverter.convertToDataType(schemaString)
val schema         = Schema.newBuilder().fromRowDataType(schemaDataType).build()

val tableDescriptor = TableDescriptor
  .forConnector("kafka")
  .schema(schema)
  .comment(topicName)
  .option(KafkaConnectorOptions.TOPIC, List(topicName).asJava)
  .option(KafkaConnectorOptions.VALUE_FORMAT, "avro-confluent")
  .option("value.avro-confluent.url", config.getString("schema.registry.url"))
  .option(KafkaConnectorOptions.PROPS_GROUP_ID, mainOutputTopicName)
  .option(KafkaConnectorOptions.PROPS_BOOTSTRAP_SERVERS, bootstrapServers)
  .build()

将Table转换为Stream的代码:

val table = tableEnv.from(tableDescriptor)

val stream = tableEnv.toDataStream(table, classOf[DataType])

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 11:28:23