beam-commits mailing list archives

Site index · List index
Message view « Date » · « Thread »
Top « Date » · « Thread »
From Lusitanian <>
Subject [GitHub] incubator-beam pull request #630: Allow custom timestamp/watermark function ...
Date Mon, 11 Jul 2016 21:11:46 GMT
GitHub user Lusitanian opened a pull request:

    Allow custom timestamp/watermark function for UnboundedFlinkSource

    So, using an `UnboundedFlinkSource` seems to force the timestamp applied to each element
of the incoming stream to the ingestion time, rather than allowing for proper event timestamping.
This PR  adds the ability to created an `UnboundedFlinkSource` with a custom TimestampAssigner,
which should alleviate that issue. Note that this is particularly useful, as currently the
only means of consuming from Kafka 0.8 using Beam/Flink Runner is to wrap Flink's Kafka 0.8

You can merge this pull request into a Git repository by running:

    $ git pull flink_event_time

Alternatively you can review and apply these changes as the patch at:

To close this pull request, make a commit to your master/trunk branch
with (at least) the following in the commit message:

    This closes #630
commit 2c5a4cd080eeaaf3505c64b773589c2da8217381
Author: David Desberg <>
Date:   2016-07-11T19:24:18Z

    Allow for custom timestamp/watermark function in FlinkPipelineRunner

commit 5d274e4ee5b5fc498bcabcab75bd7be30210bb3c
Author: David Desberg <>
Date:   2016-07-11T19:56:21Z

    Added new "of" signature and constructor for UnboundedFlinkSource to
    allow event timestamping


If your project is set up for it, you can reply to this email and have your
reply appear on GitHub as well. If your project does not have this feature
enabled and wishes so, or if the feature is enabled but not working, please
contact infrastructure at or file a JIRA ticket
with INFRA.

View raw message