This is an automated email from the ASF dual-hosted git repository.
github-bot pushed a change to tag nightly-master
in repository https://gitbox.apache.org/repos/asf/beam.git.
*** WARNING: tag nightly-master was modified! ***
from 64ec3da (commit)
to c737128 (commit)
from 64ec3da [BEAM-10892] Add Proto support to Kafka Table Provider
(#12838)
add 4836937 [BEAM-11146] Add fasterCopy option to Flink runner (#13240)
add 7507f8c [BEAM-10123] Add commit transform. (#12572)
add f67cb9a [BEAM-5504] Change Pubsub avro table jira task number in
CHANGES.md (#13248)
add 331b36e [BEAM-5570] Update javacc dependency (#13094)
add 8e093b5 [BEAM-11144] Fix trigger prefetching so that the correct
trigger index is used for the state namespace.
add 9dba074 Merge pull request #13221: [BEAM-11144] Fix trigger
prefetching so that the correct trigger index
add 092eeea [BEAM-11130] Exclude OrderedListState category for Dataflow
V2.
add 27d4fb4 Merge pull request #13203: [BEAM-11130] Exclude
OrderedListState category for Dataflow V2.
add 58b11d3 Add Dataflow Runner V2 ValidatesRunner streaming test
configuration.
add aa28f51 Merge pull request #13251: Add Dataflow Runner V2 Java
ValidatesRunner streaming test configuration.
add c3cf904 Implementing Python Bounded Source Reader DoFn (#13154)
add 0dc4afe Exclude SDF test suite because it requires support of
self-checkpoint.
add f2002c6 Merge pull request #13250 from [BEAM-11187] Exclude SDF test
suite because it requires support of self-checkpoint.
add 96982bd [BEAM-11164] Fixes bug in beam.Partition (#13236)
add c737128 [BEAM-10409] Remap all PCollections in KeyWithNone
elimination (#13204)
No new revisions were added by this update.
Summary of changes:
..._PortableValidatesRunner_Flink_Streaming.groovy | 11 +-
...tCommit_Java_ValidatesRunner_Dataflow_V2.groovy | 2 +-
...a_ValidatesRunner_Dataflow_V2_Streaming.groovy} | 8 +-
CHANGES.md | 3 +-
.../apache/beam/runners/core/ReduceFnRunner.java | 13 +-
.../core/triggers/AfterAllStateMachine.java | 21 ++
.../AfterDelayFromFirstElementStateMachine.java | 15 +-
.../core/triggers/AfterEachStateMachine.java | 21 ++
.../core/triggers/AfterFirstStateMachine.java | 21 ++
.../core/triggers/AfterPaneStateMachine.java | 19 +-
.../core/triggers/AfterWatermarkStateMachine.java | 30 ++
.../core/triggers/DefaultTriggerStateMachine.java | 9 +
.../triggers/ExecutableTriggerStateMachine.java | 12 +
.../runners/core/triggers/NeverStateMachine.java | 9 +
.../core/triggers/OrFinallyStateMachine.java | 21 ++
.../core/triggers/RepeatedlyStateMachine.java | 19 ++
.../triggers/ReshuffleTriggerStateMachine.java | 9 +
.../runners/core/triggers/TriggerStateMachine.java | 92 +++---
.../TriggerStateMachineContextFactory.java | 92 +++++-
.../core/triggers/TriggerStateMachineRunner.java | 19 +-
.../ExecutableTriggerStateMachineTest.java | 9 +
.../core/triggers/RepeatedlyStateMachineTest.java | 9 +-
.../core/triggers/StubTriggerStateMachine.java | 9 +
.../core/triggers/TriggerStateMachineTest.java | 18 ++
.../core/triggers/TriggerStateMachineTester.java | 9 +-
.../translation/types/CoderTypeSerializer.java | 31 +-
.../translation/types/CoderTypeSerializerTest.java | 6 +-
runners/flink/job-server/flink_job_server.gradle | 3 +
.../FlinkBatchPortablePipelineTranslator.java | 14 +-
.../flink/FlinkBatchTransformTranslators.java | 13 +-
.../flink/FlinkBatchTranslationContext.java | 2 +-
.../beam/runners/flink/FlinkPipelineOptions.java | 12 +
.../org/apache/beam/runners/flink/FlinkRunner.java | 5 +
.../FlinkStreamingPortablePipelineTranslator.java | 43 ++-
.../flink/FlinkStreamingTransformTranslators.java | 53 +++-
.../flink/FlinkStreamingTranslationContext.java | 2 +-
.../apache/beam/runners/flink/TestFlinkRunner.java | 3 +-
.../translation/types/CoderTypeInformation.java | 11 +-
.../wrappers/streaming/DoFnOperator.java | 35 ++-
.../streaming/ExecutableStageDoFnOperator.java | 13 +-
.../streaming/KvToByteBufferKeySelector.java | 7 +-
.../streaming/io/UnboundedSourceWrapper.java | 4 +-
.../streaming/stableinput/BufferingDoFnRunner.java | 12 +-
.../state/FlinkBroadcastStateInternals.java | 45 ++-
.../streaming/state/FlinkStateInternals.java | 95 ++++--
.../flink/FlinkExecutionEnvironmentsTest.java | 57 ++--
.../FlinkPipelineExecutionEnvironmentTest.java | 21 +-
.../runners/flink/FlinkPipelineOptionsTest.java | 10 +-
.../flink/FlinkRequiresStableInputTest.java | 3 +-
.../apache/beam/runners/flink/FlinkRunnerTest.java | 3 +-
.../beam/runners/flink/FlinkSavepointTest.java | 3 +-
.../beam/runners/flink/FlinkSubmissionTest.java | 2 +-
.../runners/flink/FlinkTransformOverridesTest.java | 3 +-
.../FlinkBroadcastStateInternalsTest.java | 7 +-
.../flink/streaming/FlinkStateInternalsTest.java | 18 +-
.../flink/streaming/GroupByWithNullValuesTest.java | 3 +-
.../wrappers/streaming/DoFnOperatorTest.java | 156 ++++++----
.../streaming/ExecutableStageDoFnOperatorTest.java | 67 +++-
.../wrappers/streaming/WindowDoFnOperatorTest.java | 9 +-
.../streaming/io/UnboundedSourceWrapperTest.java | 10 +-
.../stableinput/BufferingDoFnRunnerTest.java | 5 +-
runners/google-cloud-dataflow-java/build.gradle | 269 ++++++++++------
sdks/java/io/clickhouse/build.gradle | 2 +-
.../beam/sdk/io/kafka/KafkaCommitOffset.java | 130 ++++++++
.../java/org/apache/beam/sdk/io/kafka/KafkaIO.java | 58 +++-
.../beam/sdk/io/kafka/ReadFromKafkaDoFn.java | 13 +-
sdks/python/apache_beam/io/gcp/bigquery.py | 9 +-
sdks/python/apache_beam/io/iobase.py | 341 ++++++++++-----------
sdks/python/apache_beam/io/iobase_test.py | 58 ++--
sdks/python/apache_beam/io/textio_test.py | 14 +-
.../portability/fn_api_runner/translations.py | 17 +-
sdks/python/apache_beam/transforms/core.py | 7 +-
.../apache_beam/transforms/ptransform_test.py | 18 ++
.../shortcodes/flink_java_pipeline_options.html | 5 +
.../shortcodes/flink_python_pipeline_options.html | 5 +
75 files changed, 1510 insertions(+), 722 deletions(-)
copy
.test-infra/jenkins/{job_PostCommit_Java_ValidatesRunner_Dataflow_V2.groovy =>
job_PostCommit_Java_ValidatesRunner_Dataflow_V2_Streaming.groovy} (90%)
create mode 100644
sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaCommitOffset.java