【问题标题】:How to print an aggregated DataStream in flink?如何在 flink 中打印聚合的 DataStream?
【发布时间】: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


【解决方案1】:

使用 Flink 将集合保持在状态可能非常昂贵,因为在某些情况下,集合会经常被序列化和反序列化。如果可能,最好使用 Flink 内置的 ListState 和 MapState 类型。

这里有一个例子说明了一些事情:

public static void main(String[] args) throws Exception {

    StreamExecutionEnvironment env =
        StreamExecutionEnvironment.getExecutionEnvironment();

    env.fromElements(1L, 2L, 3L, 4L, 3L, 2L, 1L, 0L)
        .keyBy(x -> 1)
        .process(new KeyedProcessFunction<Integer, Long, List<Long>> () {
            private transient MapState<Long, Boolean> set;

            @Override
            public void open(Configuration parameters) throws Exception {
                set = getRuntimeContext().getMapState(new MapStateDescriptor<>("set", Long.class, Boolean.class));
            }

            @Override
            public void processElement(Long x, Context context, Collector<List<Long>> out) throws Exception {
                if (set.contains(x)) {
                    System.out.println("set contains " + x);
                } else {
                    set.put(x, true);
                    List<Long> list = new ArrayList<>();
                    Iterator<Long> iter = set.keys().iterator();
                    iter.forEachRemaining(list::add);
                    out.collect(list);
                }
            }
        })
        .print();

    env.execute();

}

请注意,我想使用键控状态,但事件中没有任何东西可用作键,所以我只是通过常量键控流。这通常不是一个好主意,因为它会阻止处理并行进行 - 但由于您将聚合为一个集合,因此您不能并行执行此操作,因此不会造成任何伤害。

我将 Longs 集合表示为 MapState 对象的键。当我想输出集合时,我将它收集为一个列表。当我只想打印一些东西进行调试时,我只需要使用 System.out。

我在 IDE 中运行此作业时看到的是:

[1]
[1, 2]
[1, 2, 3]
[1, 2, 3, 4]
set contains 3
set contains 2
set contains 1
[0, 1, 2, 3, 4]

如果您希望每秒查看 MapState 中的内容,可以在 process 函数中使用 Timer。

【讨论】:

  • 嗨!非常感谢,但我很难将它映射到我的模型。你能从Datastream&lt;Set&lt;Long&gt;&gt;开始吗,因为如果我有Datastream&lt;Set&lt;Long&gt;&gt;,我就不能做keyBy(x-&gt;1)这样的事情。
  • 我不明白你在做什么。请分享你拥有的东西,即使它不起作用。
  • 粘贴了引发错误的代码的sn-p。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-05-24
  • 1970-01-01
  • 1970-01-01
  • 2018-09-12
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多