【发布时间】: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