代码之家  ›  专栏  ›  技术社区  ›  user3139545

在Apache Flink测试中是否有虚拟时间的概念,就像在Reactor和RxJava中一样

  •  0
  • user3139545  · 技术社区  · 7 年前

    在RxJava和Reactor中,有一个虚拟时间的概念来测试依赖于时间的操作符。在弗林克我不知道怎么做。例如,我将下面的示例放在一起,我想在其中对迟到的事件进行处理,以了解如何处理这些事件。但是我不明白这样的测试会是什么样子?有没有办法把燧石和反应堆结合起来,使试验更好?

    public class PlayWithFlink {
    
        public static void main(String[] args) throws Exception {
    
            final OutputTag<MyEvent> lateOutputTag = new OutputTag<MyEvent>("late-data"){};
    
            // TODO understand how BoundedOutOfOrderness is related to allowedLateness
            BoundedOutOfOrdernessTimestampExtractor<MyEvent> eventTimeFunction = new BoundedOutOfOrdernessTimestampExtractor<MyEvent>(Time.seconds(10)) {
                @Override
                public long extractTimestamp(MyEvent element) {
                    return element.getEventTime();
                }
            };
    
            StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
            env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);
    
            DataStream<MyEvent> events = env.fromCollection(MyEvent.examples())
                    .assignTimestampsAndWatermarks(eventTimeFunction);
    
            AggregateFunction<MyEvent, MyAggregate, MyAggregate> aggregateFn = new AggregateFunction<MyEvent, MyAggregate, MyAggregate>() {
                @Override
                public MyAggregate createAccumulator() {
                    return new MyAggregate();
                }
    
                @Override
                public MyAggregate add(MyEvent myEvent, MyAggregate myAggregate) {
                    if (myEvent.getTracingId().equals("trace1")) {
                        myAggregate.getTrace1().add(myEvent);
                        return myAggregate;
                    }
                    myAggregate.getTrace2().add(myEvent);
                    return myAggregate;
                }
    
                @Override
                public MyAggregate getResult(MyAggregate myAggregate) {
                    return myAggregate;
                }
    
                @Override
                public MyAggregate merge(MyAggregate myAggregate, MyAggregate acc1) {
                    acc1.getTrace1().addAll(myAggregate.getTrace1());
                    acc1.getTrace2().addAll(myAggregate.getTrace2());
                    return acc1;
                }
            };
    
            KeySelector<MyEvent, String> keyFn = new KeySelector<MyEvent, String>() {
                @Override
                public String getKey(MyEvent myEvent) throws Exception {
                    return myEvent.getTracingId();
                }
            };
    
            SingleOutputStreamOperator<MyAggregate> result = events
                    .keyBy(keyFn)
                    .window(EventTimeSessionWindows.withGap(Time.seconds(10)))
                    .allowedLateness(Time.seconds(20))
                    .sideOutputLateData(lateOutputTag)
                    .aggregate(aggregateFn);
    
    
            DataStream lateStream = result.getSideOutput(lateOutputTag);
    
            result.print("SessionData");
    
            lateStream.print("LateData");
    
            env.execute();
        }
    }
    
    class MyEvent {
        private final String tracingId;
        private final Integer count;
        private final long eventTime;
    
        public MyEvent(String tracingId, Integer count, long eventTime) {
            this.tracingId = tracingId;
            this.count = count;
            this.eventTime = eventTime;
        }
    
        public String getTracingId() {
            return tracingId;
        }
    
        public Integer getCount() {
            return count;
        }
    
        public long getEventTime() {
            return eventTime;
        }
    
        public static List<MyEvent> examples() {
            long now = System.currentTimeMillis();
            MyEvent e1 = new MyEvent("trace1", 1, now);
            MyEvent e2 = new MyEvent("trace2", 1, now);
            MyEvent e3 = new MyEvent("trace2", 1, now - 1000);
            MyEvent e4 = new MyEvent("trace1", 1, now - 200);
            MyEvent e5 = new MyEvent("trace1", 1, now - 50000);
            return Arrays.asList(e1,e2,e3,e4, e5);
        }
    
        @Override
        public String toString() {
            return "MyEvent{" +
                    "tracingId='" + tracingId + '\'' +
                    ", count=" + count +
                    ", eventTime=" + eventTime +
                    '}';
        }
    }
    
    class MyAggregate {
        private final List<MyEvent> trace1 = new ArrayList<>();
        private final List<MyEvent> trace2 = new ArrayList<>();
    
    
        public List<MyEvent> getTrace1() {
            return trace1;
        }
    
        public List<MyEvent> getTrace2() {
            return trace2;
        }
    
        @Override
        public String toString() {
            return "MyAggregate{" +
                    "trace1=" + trace1 +
                    ", trace2=" + trace2 +
                    '}';
        }
    }
    

    运行this的输出是:

    SessionData:1> MyAggregate{trace1=[], trace2=[MyEvent{tracingId='trace2', count=1, eventTime=1551034666081}, MyEvent{tracingId='trace2', count=1, eventTime=1551034665081}]}
    SessionData:3> MyAggregate{trace1=[MyEvent{tracingId='trace1', count=1, eventTime=1551034166081}], trace2=[]}
    SessionData:3> MyAggregate{trace1=[MyEvent{tracingId='trace1', count=1, eventTime=1551034666081}, MyEvent{tracingId='trace1', count=1, eventTime=1551034665881}], trace2=[]}
    

    不过,我希望看到 e5

    0 回复  |  直到 7 年前
        1
  •  1
  •   David Anderson    7 年前

    如果您将水印分配程序修改为这样

    AssignerWithPunctuatedWatermarks eventTimeFunction = new AssignerWithPunctuatedWatermarks<MyEvent>() {
        long maxTs = 0;
    
        @Override
        public long extractTimestamp(MyEvent myEvent, long l) {
            long ts = myEvent.getEventTime();
            if (ts > maxTs) {
                maxTs = ts;
            }
            return ts;
        }
    
        @Override
        public Watermark checkAndGetNextWatermark(MyEvent event, long extractedTimestamp) {
            return new Watermark(maxTs - 10000);
        }
    };
    

    然后你就会得到你期望的结果。我不是推荐这个——只是用它来说明发生了什么。

    BoundedOutOfOrdernessTimestampExtractor 是一个周期性的水印生成器,它只会每隔200毫秒向流中插入一个水印(默认情况下)。因为你的工作早在那之前就完成了,所以你的工作所经历的唯一的水印就是Flink在每个有限流的末尾注入的水印(带有MAX_值的水印)。迟到与水印有关,而您预期迟到的事件正设法在水印之前到达。

    通过切换到标点水印,您可以强制水印更频繁地出现,或更精确地出现在流中的特定点。这通常是不必要的(而且过频繁的水印会导致开销),但是当您想要对水印的序列有很强的控制时,这是很有帮助的。

    至于如何编写测试,您可以查看 test harnesses 用于弗林克自己的测试,或者 flink-spector

    更新:

    与boundedOutfordernessTimestampExtractor关联的时间间隔是流的无序程度的规范。在此范围内到达的事件不会被视为延迟,并且事件时间计时器在该延迟过去之前不会触发,从而为无序事件提供到达时间。allowedLateness只适用于window API,并描述框架保持窗口状态的时间超过正常窗口触发时间的长度,以便事件仍然可以添加到窗口并导致延迟触发。在这个额外的间隔之后,窗口状态被清除,随后的事件被发送到侧输出(如果配置)。

    enter image description here

    所以当你使用 BoundedOutOfOrdernessTimestampExtractor<MyEvent>(Time.seconds(10)) 你是 最多 10秒以防更早的事件到来。(如果您正在处理历史数据,那么您可以在1秒内处理10秒的数据,或者不处理——知道您将等待n秒的事件时间过去,这并不能说明实际需要多长时间。)

    有关此主题的详细信息,请参见 Event Time and Watermarks

    推荐文章