This is an automated email from the ASF dual-hosted git repository.
github-bot pushed a change to branch nightly-refs/heads/master
in repository https://gitbox.apache.org/repos/asf/beam.git
from bfa0c59ebcd The Bag Partition is now configurable. (#33805)
add b701737261e increase timeout (#33835)
add e380edf1974 Change runner to ubuntu-22.04 (#33816)
add 83b373570fe [flink-runner] Improve Datastream for batch performances
(#32440)
add d63b8fb760e Fix Playground CI Nightly job (#33837)
add df13ffe96d6 sdks/python: enable named aggregation to deferred
DataFrame groupby (#33672)
add c7b2695f364 Adding handling of Nulled lists to beam_row_from_dict
(#33830)
add edf7c9067cc fix activate (#33839)
add 8cbccc9145d Kafka source offset-based deduplication. (#33596)
add 4d7802877ea [Prism] Disable trie tests. (#33843)
add 0230b5d420f Update Dataflow api client version (#33829)
add 5aae10db7af Add a flag to disable BoundedTrie metrics in Beam
add 8805ab2240c Merge pull request #33457 from
rohitsinha54/btrie-disable-flag
No new revisions were added by this update.
Summary of changes:
.../trigger_files/beam_PostCommit_Go_VR_Flink.json | 1 +
.../beam_PostCommit_Java_Examples_Flink.json | 3 +-
.../beam_PostCommit_Java_PVR_Flink_Batch.json | 1 +
.../beam_PostCommit_Java_PVR_Flink_Docker.json | 2 +
.../beam_PostCommit_Java_PVR_Flink_Streaming.json | 3 +-
...beam_PostCommit_Java_ValidatesRunner_Flink.json | 1 +
.github/trigger_files/beam_PostCommit_Python.json | 1 +
.../trigger_files/beam_PostCommit_XVR_Flink.json | 3 +-
.../workflows/beam_Publish_Beam_SDK_Snapshots.yml | 2 +-
.../org/apache/beam/gradle/BeamModulePlugin.groovy | 2 +-
.../katas/coretransforms/flattenWith/Task.java | 2 +-
.../learning/katas/coretransforms/tee/Task.java | 13 +
.../FlattenWith/FlattenWith/task.py | 28 +-
.../katas/python/Core Transforms/Tee/Tee/task.py | 24 +-
playground/categories.yaml | 1 +
.../core/GroupAlsoByWindowViaWindowSetNewDoFn.java | 6 +-
.../translation/types/CoderTypeSerializer.java | 24 +-
runners/flink/flink_runner.gradle | 12 +-
.../runners/flink/FlinkExecutionEnvironments.java | 10 +
.../beam/runners/flink/FlinkPipelineOptions.java | 9 +-
.../FlinkStreamingAggregationsTranslators.java | 544 +++++++++++++++++++++
.../flink/FlinkStreamingPipelineTranslator.java | 5 +-
.../FlinkStreamingPortablePipelineTranslator.java | 35 +-
.../flink/FlinkStreamingTransformTranslators.java | 441 +++++------------
.../beam/runners/flink/adapter/FlinkKey.java | 87 ++++
.../wrappers/streaming/DoFnOperator.java | 118 +++--
.../streaming/ExecutableStageDoFnOperator.java | 81 +--
...ySelector.java => KvToFlinkKeyKeySelector.java} | 23 +-
.../streaming/PartialReduceBundleOperator.java | 181 +++++++
...eySelector.java => SdfFlinkKeyKeySelector.java} | 27 +-
.../wrappers/streaming/SplittableDoFnOperator.java | 5 +-
.../wrappers/streaming/WindowDoFnOperator.java | 22 +-
.../wrappers/streaming/WorkItemKeySelector.java | 20 +-
.../wrappers/streaming/io/DedupingOperator.java | 8 +-
.../wrappers/streaming/io/source/FlinkSource.java | 27 +-
...or.java => LazyFlinkSourceSplitEnumerator.java} | 145 +++---
.../source/bounded/FlinkBoundedSourceReader.java | 6 +
.../streaming/state/FlinkStateInternals.java | 450 +++++++++++------
.../runners/flink/FlinkPipelineOptionsTest.java | 6 +-
.../beam/runners/flink/FlinkSubmissionTest.java | 48 +-
.../beam/runners/flink/adapter/FlinkKeyTest.java | 94 ++++
.../flink/streaming/FlinkStateInternalsTest.java | 30 +-
.../wrappers/streaming/DoFnOperatorTest.java | 206 ++++----
.../streaming/ExecutableStageDoFnOperatorTest.java | 34 +-
.../wrappers/streaming/WindowDoFnOperatorTest.java | 49 +-
.../flink-conf.yaml | 6 +-
runners/prism/java/build.gradle | 3 +
.../container/license_scripts/license_script.sh | 8 +-
.../java/org/apache/beam/sdk/metrics/Metrics.java | 21 +-
.../org/apache/beam/sdk/metrics/MetricsTest.java | 7 +
.../beam/sdk/io/kafka/KafkaCheckpointMark.java | 21 +
.../java/org/apache/beam/sdk/io/kafka/KafkaIO.java | 53 +-
.../KafkaIOReadImplementationCompatibility.java | 1 +
.../org/apache/beam/sdk/io/kafka/KafkaIOUtils.java | 13 +
.../beam/sdk/io/kafka/KafkaUnboundedReader.java | 33 +-
.../beam/sdk/io/kafka/KafkaUnboundedSource.java | 23 +-
.../beam/sdk/io/kafka/KafkaIOExternalTest.java | 1 -
...KafkaIOReadImplementationCompatibilityTest.java | 15 +-
.../org/apache/beam/sdk/io/kafka/KafkaIOTest.java | 61 ++-
.../sdk/io/kafka/upgrade/KafkaIOTranslation.java | 10 +
.../io/kafka/upgrade/KafkaIOTranslationTest.java | 1 +
sdks/python/apache_beam/dataframe/frames.py | 95 +++-
sdks/python/apache_beam/dataframe/frames_test.py | 10 +
sdks/python/apache_beam/io/gcp/bigquery_tools.py | 6 +-
.../apache_beam/io/gcp/bigquery_tools_test.py | 11 +-
sdks/python/container/run_validatescontainer.sh | 2 +-
.../shortcodes/flink_java_pipeline_options.html | 5 +
.../shortcodes/flink_python_pipeline_options.html | 5 +
68 files changed, 2248 insertions(+), 1003 deletions(-)
create mode 100644
runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkStreamingAggregationsTranslators.java
create mode 100644
runners/flink/src/main/java/org/apache/beam/runners/flink/adapter/FlinkKey.java
rename
runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/{KvToByteBufferKeySelector.java
=> KvToFlinkKeyKeySelector.java} (62%)
create mode 100644
runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/PartialReduceBundleOperator.java
rename
runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/{SdfByteBufferKeySelector.java
=> SdfFlinkKeyKeySelector.java} (65%)
copy
runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/source/{FlinkSourceSplitEnumerator.java
=> LazyFlinkSourceSplitEnumerator.java} (55%)
create mode 100644
runners/flink/src/test/java/org/apache/beam/runners/flink/adapter/FlinkKeyTest.java
copy runners/flink/src/test/{resources =>
validatesRunnerConfig}/flink-conf.yaml (88%)