flink-user mailing list archives

Site index · List index
Message view « Date » · « Thread »
Top « Date » · « Thread »
From Anil <anilsingh....@gmail.com>
Subject Flink Checkpoint Barrier in case of Join
Date Tue, 02 Oct 2018 21:31:43 GMT
I'm trying to understand when will Flink's Stream Barrier (for checkpoint) be
emitted by the join operator. 

Consider a query like -

select * from  stream_1 a1 INNER JOIN  stream_2 a2 on a2.orderId =
a1.orderId group by HOP(a1.proctime, INTERVAL '1' HOUR, INTERVAL '1' DAY),

Since I'm using a Hopping window on 1 day here, Flink will have to cache my
entire 1 day events. 
The join operator will receive stream barrier from the previous operator. 
Join operator will emit one stream barrier but I'm not sure on what basis
and when will it be emitted. 

Any help will be appreciated. Thanks!

>From Flink's documentation - 

```A core element in Flink’s distributed snapshotting are the stream
barriers. These barriers are injected into the data stream and flow with the
records as part of the data stream.  The point where the barriers for
snapshot n are injected (let’s call it Sn) is the position in the source
stream up to which the snapshot covers the data. The barriers then flow
downstream. When an intermediate operator has received a barrier for
snapshot n from all of its input streams, it emits a barrier for snapshot n
into all of its outgoing streams. Once a sink operator (the end of a
streaming DAG) has received the barrier n from all of its input streams, it
acknowledges that snapshot n to the checkpoint coordinator. After all sinks
have acknowledged a snapshot, it is considered completed.```

Sent from: http://apache-flink-user-mailing-list-archive.2336050.n4.nabble.com/

View raw message