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

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的泛型参数,从根源解决类型推断问题:

  1. 初始化KafkaSource.builder时,显式声明泛型为你的事件类CreatedEvent
  2. 替换反序列化配置为setValueOnlyDeserializer,传入你自定义的CreatedEventSchema,仅解析消息value
  3. 移除流模式下不需要的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 22:21:32