samza-dev mailing list archives

Site index · List index
Message view « Date » · « Thread »
Top « Date » · « Thread »
From Jagadish Venkatraman <>
Subject Re: Review Request 51346: SAMZA-974 - Support finite datasources in Samza that have a notion of End-Of-Stream
Date Tue, 06 Sep 2016 16:02:05 GMT

This is an automatically generated e-mail. To reply, visit:

(Updated Sept. 6, 2016, 4:02 p.m.)

Review request for samza, Boris Shkolnik, Chris Pettitt, Yi Pan (Data Infrastructure), Navina
Ramesh, and Xinyu Liu.

Repository: samza


Samza currently works with unbounded data sources (kafka streams). However, for bounded data
sources like HDFS files, snapshot files which are not infinite, we need a notion of 'end-of-stream'.

This is a step towards realizing a 'finite' Samza job that terminates once data processing
is complete.(as opposed to an infinite stream job that keeps running)  

RB changes:
- New interface for EndOfStreamListener
- New 'end-of-stream' state in the state-machine of AsyncStreamTask (Invariant: When end-of-stream
is reached there are no buffered messages, no-callbacks are in-flight and no-window or commit
call shall be in progress)
- Changes to allow clean shut-downs of the tasks/container/job for end-of-stream.

Design Doc and Implementation Notes:

Diffs (updated)

  build.gradle 1d4eb74b1294318db8454631ddd0901596121ab2 
  checkstyle/import-control.xml 7e77702bcd5c32f7fdaf1558337505993c1abe06 
  samza-api/src/main/java/org/apache/samza/system/ cc860cf7eb4d514736913c1dceaa80534b61d71a

  samza-api/src/main/java/org/apache/samza/system/ a8f858aa7e4f4ce436f450cf439fe1a102983c64

  samza-api/src/main/java/org/apache/samza/task/ PRE-CREATION

  samza-core/src/main/java/org/apache/samza/task/ a510bb0c5914c772438930d27f100b4d360c1296

  samza-core/src/main/java/org/apache/samza/task/ 1fc645673b7547b642830df5639c0b4fcd11c0d5

  samza-core/src/main/scala/org/apache/samza/container/TaskInstance.scala 89f6857014489aba2db4129bc2e26dfec5b10652

  samza-core/src/main/scala/org/apache/samza/system/SystemConsumers.scala a8355b944cad54faacf5eeb883d8f4b630440757

  samza-core/src/test/java/org/apache/samza/task/ ca913dea79fecbcecdfd1010dc794318055c5764



Unit tests to test scenarios for inorder processing, out-of-order processing and commit semantics.


Jagadish Venkatraman

  • Unnamed multipart/alternative (inline, None, 0 bytes)
View raw message