【发布时间】:2020-03-05 02:16:50
【问题描述】:
我有一个以Set<Long> 表示的自定义状态计算,它会随着我的Datastream<Set<Long>> 看到来自Kafka 的新事件而不断更新。现在,每次更新我的状态时,我都想将更新后的状态打印到标准输出。想知道如何在 Flink 中做到这一点?对所有窗口和触发器操作有点困惑,我不断收到以下错误。
Caused by: java.lang.RuntimeException: Record has Long.MIN_VALUE timestamp (= no timestamp marker). Is the time characteristic set to 'ProcessingTime', or did you forget to call 'DataStream.assignTimestampsAndWatermarks(...)'?
我只想知道如何将我的聚合流Datastream<Set<Long>> 打印到标准输出或写回另一个 kafka 主题?
下面是引发错误的代码的sn-p。
StreamTableEnvironment bsTableEnv = StreamTableEnvironment.create(env, bsSettings);
DataStream<Set<Long>> stream = bsTableEnv.toAppendStream(kafkaSourceTable, Row.class)
stream
.aggregate(new MyCustomAggregation(100))
.process(new ProcessFunction<Set<Long>, Object>() {
@Override
public void processElement(Set<Long> value, Context ctx, Collector<Object> out) throws Exception {
System.out.println(value.toString());
}
});
【问题讨论】:
-
请详细说明您想要完成的任务。在每个事件之后输出整个 Set 会很昂贵,特别是如果 Set 随着每个事件而增长。这是为了调试,还是???
-
是的,你明白了!它主要用于调试,所以每秒输出对我也有好处!不需要在每个事件之后输出。
-
我无法从您共享的代码中弄清楚发生了什么。关于时间戳和水印的错误只是从 Flink 的窗口代码中抛出的,我没有看到任何窗口。此外,DataStreams 上没有聚合方法——仅在 Windows 上。我认为打印可能会起作用,但在到达那里之前工作就失败了。
标签: apache-flink