Flink调用fromSource创建KafkaSource时类型为Nothing如何解决
问题原因
编译时流元素被识别为Nothing类型,核心是Scala泛型自动推断失效:
- 调用
KafkaSource.builder时未显式指定泛型参数,编译器默认将builder的泛型判定为Nothing,最终build()得到的Kafka Source自然也携带Nothing类型 - 原代码使用的
KafkaRecordDeserializationSchema.of(new CreatedEventSchema)仅负责反序列化逻辑,无法为编译器提供足够的泛型信息,无法自动推导出流元素类型为CreatedEvent - 原代码还存在两处逻辑问题:一是同时配置了
setBounded(OffsetsInitializer.latest)(批模式停止偏移配置)和setStartingOffsets(OffsetsInitializer.earliest),流模式下二者冲突;二是JDBC写入逻辑中第二个参数的索引误写为1,会覆盖第一个timestamp字段的值。
修复方案
因为你仅需要处理Kafka消息的value部分,不需要解析key、分区、偏移量等元数据,直接使用Kafka Source提供的value专用反序列化配置即可,同时显式指定builder的泛型参数,从根源解决类型推断问题:
- 初始化
KafkaSource.builder时,显式声明泛型为你的事件类CreatedEvent - 替换反序列化配置为
setValueOnlyDeserializer,传入你自定义的CreatedEventSchema,仅解析消息value - 移除流模式下不需要的
setBounded配置,修正JDBC参数索引错误
修正后的可编译代码如下:
import org.apache.flink.api.common.eventtime.WatermarkStrategy import org.apache.flink.connector.jdbc.{JdbcConnectionOptions, JdbcExecutionOptions, JdbcSink} import org.apache.flink.connector.kafka.source.KafkaSource import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer // 其余导入根据自身项目补充,包括CreatedEvent、CreatedEventSchema等 object Main { def main(args: Array[String]) { // 显式指定泛型为CreatedEvent,避免类型推断为Nothing val kafkaSource = KafkaSource.builder[CreatedEvent]() .setBootstrapServers("localhost:29092") .setProperty("partition.discovery.interval.ms", "10000") .setTopics("created") .setStartingOffsets(OffsetsInitializer.earliest) // 仅解析消息value,传入自定义反序列化schema .setValueOnlyDeserializer(new CreatedEventSchema) .build() val env = StreamExecutionEnvironment.getExecutionEnvironment val stream = env.fromSource(kafkaSource, WatermarkStrategy.noWatermarks(), "Kafka Source") stream.addSink(JdbcSink.sink( "INSERT INTO conversations (timestamp, active_conversations, total_conversations) VALUES (?,?,?)", (statement, event) => { statement.setTime(1, event.date) statement.setInt(2, event.a) // 修正原代码索引错误,第二个参数对应JDBC索引2 statement.setInt(3, event.b) }, JdbcExecutionOptions.builder() .withBatchSize(1000) .withBatchIntervalMs(200) .withMaxRetries(5) .build(), new JdbcConnectionOptions.JdbcConnectionOptionsBuilder() .withUrl("jdbc:postgresql://localhost:5432/reporting") .withDriverName("org.postgresql.Driver") .withUsername("postgres") .withPassword("veryverysecret:-)") .build() )) env.execute() } }
如果你需要保留KafkaRecordDeserializationSchema以获取消息元数据,不需要替换反序列化器,只要在初始化builder时加上[CreatedEvent]泛型声明,同样可以解决Nothing类型的编译问题,只是需要在后续逻辑中自行从反序列化结果中提取消息value。
内容的提问来源于stack exchange,提问作者Sjoerd222888
相关产品推荐
相关产品推荐

