Scala-如何过滤KStream(Kafka Streams)



我是Scala的新手,我正在尝试根据第二个组件字段过滤KStream[String,JsonNode]。

例如,工作的Java代码是这样的:

import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.kstream.Consumed;
import org.apache.kafka.streams.kstream.KStream;
import org.apache.kafka.streams.kstream.Produced;
...
import com.fasterxml.jackson.databind.JsonNode;
...
...
final KStream<String, JsonNode> source = streamsBuilder.stream(inputTopic,
Consumed.with(Serdes.String(), jsonSerde));

// filter and producer preprocessed
source.filter((k, v) -> v.get("total_cost").asDouble() > 0 && v.get("num_items").asInt() > 0)
.to(outputTopic, Produced.with(Serdes.String(), jsonSerde));

我试过这个:

import org.apache.kafka.streams.kstream.{Produced,Consumed,KStream};
import org.apache.kafka.streams.StreamsBuilder;
...
import com.fasterxml.jackson.databind.JsonNode;
...
var source:KStream[String, JsonNode] = streamsBuilder.stream(inputTopic, Consumed.`with`(Serdes.String(), jsonSerde));
source.filter({
case (k:String ,v:JsonNode) => 
(v.get("total_cost").asDouble() > 0 && v.get("num_items").asInt() > 0)
})
.to(outputTopic, Produced.`with`(Serdes.String(), jsonSerde));

在上面的尝试中,我得到了:

描述资源路径位置类型缺少的参数类型展开函数匿名函数的参数类型必须是众所周知。(SLS 8.5(预期类型为:org.apache.kafka.stream.kstream.Predicate[?>:字符串,?>:com.fasterxml.jackson.databind.JsonNode]

我也尝试过这个:

source.filter((_._2.get("total_cost").asDouble() > 0 && _._2.get("num_items").asInt() > 0))
.to(outputTopic, Produced.`with`(Serdes.String(), jsonSerde));

如何在Scala中过滤这个对象?提前谢谢。

org.apache.kafka.streams.kstream.Predicate来自JavaAPI,在Scala2.11(我假设你使用它(中,你必须显式地实现接口,因此:

source.filter(new Predicate[String, JsonNode]() {
override def test(k: String, v: JsonNode): Boolean = {
v.get("total_cost").asDouble() > 0 && v.get("num_items").asInt() > 0
}
}).to(outputTopic, Produced.`with`(Serdes.String(), jsonSerde));

应该起作用。

关于SAMs(单一抽象方法(的更多信息,请点击

请注意,您不必使用Java API——这里有一流的Scala API。

相关内容

  • 没有找到相关文章

最新更新