> 在 2018年7月31日,15:47,vino yang <[email protected]> 写道:
> 
> Hi anna,
> 
> 1. The srcstream is a very high volume stream and the window size is 2 weeks 
> and 4 weeks. Is the window size a problem? In this case, I think it is not a 
> problem because I am using reduce which stores only 1 value per window. Is 
> that right?
> 
> >> Window Size is based on your business needs settings. However, if the 
> >> window size is too large, the status of the job will be large, which will 
> >> result in a longer recovery failure. You need to be aware of this. One 
> >> value per window is just a value calculated by the window. It caches all 
> >> data for the period of time before the window is triggered.
> 
> 2. I am having 2 output operations one with 2 weeks window and the other with 
> 4 weeks window. Are they executed in parallel or in sequence?
> 
> >> These two windows are calculated in parallel.
> 
> 3. When I have multiple output operations like in this case should I break it 
> into 2 different jobs ?
> 
> >> Both modes are ok. When there is only one job, the two windows will share 
> >> the source stream, but this will result in a larger state of the job and a 
> >> slower recovery. When split into two jobs, there will be two consumptions 
> >> of kafka, but the two windows are independent in both jobs.
> 
> 4. Can I run multiple jobs on the same cluster?
> 
> >> For Standalone cluster mode or Yarn Flink Session mode, etc., there is no 
> >> problem. For Flink on yarn single job mode, a cluster can usually only run 
> >> one job, which is the recommended mode.
> 
> Thanks, vino.
> 
> 2018-07-31 15:11 GMT+08:00 anna stax <[email protected]>:
> Hi all,
> 
> I am not sure when I should go for multiple jobs or have 1 job with all the 
> sources and sinks. Following is my code.
> 
>    val env = StreamExecutionEnvironment.getExecutionEnvironment
>     .......
>     // create a Kafka source
>     val srcstream = env.addSource(consumer)
> 
>     srcstream
>       .keyBy(0)
>       .window(ProcessingTimeSessionWindows.withGap(Time.days(14)))
>       .reduce  ...
>       .map ...
>       .addSink ...  
> 
>     srcstream
>       .keyBy(0)
>       .window(ProcessingTimeSessionWindows.withGap(Time.days(28)))
>       .reduce  ...
>       .map ...
>       .addSink ...
> 
>     env.execute("Job1")
> 
> My questions
> 
> 1. The srcstream is a very high volume stream and the window size is 2 weeks 
> and 4 weeks. Is the window size a problem? In this case, I think it is not a 
> problem because I am using reduce which stores only 1 value per window. Is 
> that right?
> 
> 2. I am having 2 output operations one with 2 weeks window and the other with 
> 4 weeks window. Are they executed in parallel or in sequence?
> 
> 3. When I have multiple output operations like in this case should I break it 
> into 2 different jobs ?
> 
> 4. Can I run multiple jobs on the same cluster?
> 
> Thanks
> 
>  
> 

Hi yang, 

I’m a bit confused about the window data cache that you mentioned.

> It caches all data for the period of time before the window is triggered.

In my understanding, window functions process elements incrementally unless the 
low level API ProcessWindowFunction was used, so caching data should not be 
required in most scenarios. Would you mind giving more details of the window 
caching design? And please correct me if I’m wrong. Thanks a lot.

Best regards, 
Paul Lam

Reply via email to