Hello David,

Your observation is correct, the stream time reasoning is dependent on the
buffered records from each of the input topic-partitions, and hence is
"data-driven".

Currently to get around this I'd recommend letting the producer to send
certain "marker" messages periodically to ensure stream time advance; and
for the near future we are working on improving our stream time reasoning
logic, and hopefully resolving such data-driven dependencies in the next
release:

https://issues.apache.org/jira/browse/KAFKA-3514


Guozhang


On Sat, Oct 8, 2016 at 12:29 PM, David Garcia <[email protected]> wrote:

> Actually, I think the bug is more subtle.  What happens when a consumed
> topic stops receiving messages?  The smallest timestamp will always be the
> static timestamp of this topic.
>
> -David
>
> On 10/7/16, 5:03 PM, "David Garcia" <[email protected]> wrote:
>
>     Ok I found the bug.  Basically, if there is an empty topic (in the
> list of topics being consumed), any partition-group with partitions from
> the topic will always return -1 as the smallest timestamp (see
> PartitionGroup.java).
>
>     To reproduce, simply start a kstreams consumer with one or more empty
> topics.  Punctuate will never be called.
>
>     -David
>
>     On 10/7/16, 1:11 PM, "David Garcia" <[email protected]> wrote:
>
>         Yeah, this is possible.  We have run the application (and have
> confirmed data is being received) for over 30 mins…with a 60-second timer.
> So, do we need to just rebuild our cluster with bigger machines?
>
>         -David
>
>         On 10/7/16, 11:18 AM, "Michael Noll" <[email protected]> wrote:
>
>             David,
>
>             punctuate() is still data-driven at this point, even when
> you're using the
>             WallClock timestamp extractor.
>
>             To use an example: Imagine you have configured punctuate() to
> be run every
>             5 seconds.  If there's no data being received for a minute,
> then punctuate
>             won't be called -- even though you probably would have
> expected this to
>             happen 12 times during this 1 minute.
>
>             (FWIW, there's an ongoing discussion to improve punctuate(),
> part of which
>             is motivated by the current behavior that arguably is not very
> intuitive to
>             many users.)
>
>             Could this be the problem you're seeing?  See also the related
> discussion
>             at
>             http://stackoverflow.com/questions/39535201/kafka-
> problems-with-timestampextractor
>             .
>
>
>
>
>
>
>             On Fri, Oct 7, 2016 at 6:07 PM, David Garcia <
> [email protected]> wrote:
>
>             > Hello, I’m sure this question has been asked many times.
>             > We have a test-cluster (confluent 3.0.0 release) of 3 aws
> m4.xlarges.  We
>             > have an application that needs to use the punctuate()
> function to do some
>             > work on a regular interval.  We are using the WallClock
> extractor.
>             > Unfortunately, the method is never called.  I have checked
> the
>             > filedescriptor setting for both the user as well as the
> process, and
>             > everything seems to be fine.  Is this a known bug, or is
> there something
>             > obvious I’m missing?
>             >
>             > One note, the application used to work on this cluster, but
> now it’s not
>             > working.  Not really sure what is going on?
>             >
>             > -David
>             >
>
>
>
>
>
>
>


-- 
-- Guozhang

Reply via email to