【发布时间】:2018-05-18 13:47:21
【问题描述】:
我有 2 个名为“警报”和“干预”的流,其中包含 JSON。如果连接了警报和干预,则它们将具有相同的键。我想联系他们以检测所有在 24 小时前未进行干预的警报。
但是这个程序不起作用,结果给了我所有的警报,就好像 24 小时前没有进行任何干预一样。
我重新检查了我的数据集 5 次,并且有些警报在警报日期前不到 24 小时内完成了干预。
这张图说明情况:
enter image description here
所以我需要知道警报前是否有干预。
程序代码:
final KStream<String, JsonNode> alarm = ...;
final KStream<String, JsonNode> intervention = ...;
final JoinWindows jw = JoinWindows.of(TimeUnit.HOURS.toMillis(24)).before(TimeUnit.HOURS.toMillis(24)).after(0);
final KStream<String, JsonNode> joinedAI = alarm.filter((String key, JsonNode value) -> {
return value != null;
}).leftJoin(intervention, (JsonNode leftValue, JsonNode rightValue) -> {
ObjectMapper mapper = new ObjectMapper();
JsonNode actualObj = null;
if (rightValue == null) {//No intervention before
try {
actualObj = mapper.readTree("{\"date\":\"" + leftValue.get("date").asText() + "\","
+ "\"alarm\":" + leftValue.toString()
+ "}");
} catch (IOException ex) {
Logger.getLogger(Main.class.getName()).log(Level.SEVERE, null, ex);
}
return actualObj;
} else {
return null;
}
}, jw, Joined.with(Serdes.String(), jsonSerde, jsonSerde));
final KStream<String, JsonNode> fraude = joinedAI.filter((String key, JsonNode value) -> {
return value != null;
});
fraude.foreach((key, value) -> {
rl.println("Fraude=" + key + " => " + value);
System.out.println("Fraude=" + key + " => " + value);
});
final KafkaStreams streams = new KafkaStreams(builder.build(), streamingConfig);
streams.cleanUp();
streams.start();
Runtime.getRuntime().addShutdownHook(new Thread(new Runnable() {
@Override
public void run() {
streams.close();
rl.close();
el.close();
nfl.close();
}
}));
综上所述,我想检测红色矩形enter image description here中的图案
P.S:我确保在报警记录之前发送干预记录
【问题讨论】:
-
这篇博文可能会有所帮助:confluent.io/blog/crossing-streams-joins-apache-kafka -- Kafka Streams 中的连接与 SQL 连接的语义略有不同。
标签: apache-kafka left-join apache-kafka-streams