【问题标题】:KafkaUtils API | offset management | Spark StreamingKafkaUtils API |偏移管理|火花流
【发布时间】:2016-09-13 14:20:23
【问题描述】:

我正在尝试为仅一次语义管理 kafka 偏移量。

使用偏移图创建直接流时面临的问题如下:

val fromOffsets : (TopicAndPartition, Long) = TopicAndPartition(metrics_rs.getString(1), metrics_rs.getInt(2)) -> metrics_rs.getLong(3)

KafkaUtils.createDirectStream[String, String, StringDecoder, StringDecoder,(String, String)] (ssc,kafkaParams,fromOffsets,messageHandler)

这里,

val messageHandler =
      (mmd: MessageAndMetadata[String, String]) => mmd.message.length

还有

metrics_rs = metricsStatement.executeQuery("SELECT part,off from metrics.txn_offsets where topic='"+t+''' )

我想我在声明风格上做错了......如果你能帮忙的话。 编译错误说“createDirectStream 的类型参数太多”

【问题讨论】:

  • 你知道最新的 Kafka 0.10+ 兼容的KafkaUtils.createDirectStream吗?不知道为什么要用5-type 0.8兼容的接口?

标签: scala apache-spark apache-kafka spark-streaming


【解决方案1】:

我发现你做错了几件事。

您需要传递一个Map[TopicAndPartition, Long],而目前您有一个Tuple2[TopicAndPartition, Long]。所以你需要:

val fromOffsets: Map[TopicAndPartition, Long] = 
    Map(TopicAndPartition(metrics_rs.getString(1), 
                          metrics_rs.getInt(2)) -> metrics_rs.getLong(3))

你说你来自createDirectStream 的返回类型是一个(String, String) 类型的元组,但你的messageHandler 值是一个Int。如果你想返回一个带有键值对的元组,你需要:

val messageHandler: MessageAndMetadata[String, String] => (String, String) =
  (mmd: MessageAndMetadata[String, String]) => (mmd.key(), mmd.message())

修复后,这应该可以编译:

val stream = KafkaUtils
              .createDirectStream[String, String,
                      StringDecoder, StringDecoder,
                      (String, String)] (ssc, 
                                         kafkaParams, 
                                         fromOffsets, 
                                         messageHandler)

【讨论】:

  • 是的,我确实将元组转换为 Map 和 messageHandler 但它仍然没有编译。
  • kafkaParamsSet[String] 吗?
  • 这是一个地图... val kafkaParams = Map[String, String]( "zookeeper.connect" -> KafkaConfig.zookeeperHost, "group.id" -> KafkaConfig.groupId, "auto. commit.enable" -> KafkaConfig.autoCommitEnabled, "auto.commit.interval.ms" -> KafkaConfig.autoCommitInterval, "bootstrap.servers" -> KafkaConfig.bootstrapServers ...)
  • @user1521672 是的,我的错。你确定你正确地修复了你的messageHandler 吗?错误信息是什么意思?
  • 是的... val messageHandler: (String, String) = (mmd: MessageAndMetadata[String, String]) => (mmd.key(), mmd.message()) 同样的错误“ createDirectStream 的类型参数太多"
猜你喜欢
  • 2017-01-12
  • 2016-12-19
  • 2017-04-27
  • 2021-11-16
  • 2018-08-09
  • 1970-01-01
  • 2016-06-25
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多