Flink Table API读取某Kafka Topic时列值错乱问题排查求助
Flink Table API转DataStream时Avro列值错乱问题
我用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
相关产品推荐
相关产品推荐

