This is an automated email from the ASF dual-hosted git repository.
github-actions[bot] pushed a change to branch nightly-refs/heads/master
in repository https://gitbox.apache.org/repos/asf/beam.git
from 790468a6591 Ensure 0 backlog is sent when finishing processing
restrictions. (#40102)
add a0adae53810 Kafka Streams runner: add a flush marker payload variant
add 60205a1a12a Carry target partitions in the flush marker, and shorten
the comments
add 07c1a1350b8 Merge pull request #40068: Kafka Streams runner: add a
flush marker payload variant
add 90ecd1600e3 Bump go.mongodb.org/mongo-driver from 1.17.9 to 1.17.10 in
/sdks (#40111)
add 7bf5506f9a2 Bump github/codeql-action from 4.37.9 to 4.38.0 (#40115)
add fbe8ddf6c64 Bump zizmorcore/zizmor-action from 0.6.3 to 0.6.4 (#40117)
add 898b70e3061 Bump cloud.google.com/go/storage from 1.67.0 to 1.67.1 in
/sdks (#40118)
add 830d7b63ab5 [Spark] Make cancel() cancel the Spark jobs and stop only
a session the runner created (#40103)
add 8f02ad3868d AddFiles: CommitSchemaUnion (#40104)
add 544e40f5a89 [Iceberg CDC sink] Split late data and create commit
windows (#40006)
add f9064240c35 [Iceberg CDC sink] table setup (#40007)
add a7b24f3d577 Revert "Revert "[Experimental] Use zstd compression in
Docker ...
No new revisions were added by this update.
Summary of changes:
...eam_PostCommit_Java_ValidatesRunner_Spark4.json | 2 +-
...a_ValidatesRunner_SparkStructuredStreaming.json | 3 +-
...tCommit_Python_ValidatesContainer_Dataflow.json | 2 +-
.../beam_PreCommit_Flink_Container.json | 2 +-
.github/workflows/beam_PreCommit_GHA.yml | 2 +-
.github/workflows/codeql.yml | 4 +-
.../org/apache/beam/gradle/BeamDockerPlugin.groovy | 24 +-
.../src/main/proto/kafka_streams_payload.proto | 12 +-
.../{package-info.java => FlushPayload.java} | 21 +-
.../kafka/streams/translation/KStreamsPayload.java | 73 +-
.../streams/translation/KStreamsPayloadSerde.java | 14 +-
.../translation/KStreamsPayloadSerdeTest.java | 28 +
.../translation/StreamingEvaluationContext.java | 11 +-
.../SparkStructuredStreamingPipelineResult.java | 74 +-
.../SparkStructuredStreamingRunner.java | 48 +-
.../translation/EvaluationContext.java | 20 +-
.../translation/SparkSessionFactory.java | 48 +-
...SparkStructuredStreamingPipelineResultTest.java | 66 ++
.../StructuredStreamingPipelineStateTest.java | 72 ++
sdks/go.mod | 4 +-
sdks/go.sum | 8 +-
sdks/java/io/iceberg/build.gradle | 1 +
.../beam/sdk/io/iceberg/CommitSchemaUnion.java | 402 ++++++++
.../beam/sdk/io/iceberg/cdc/sink/CommitToken.java | 314 ++++++
.../sdk/io/iceberg/cdc/sink/CommitWindows.java | 226 +++++
.../sink/DestinationShard.java} | 50 +-
.../io/iceberg/cdc/sink/PartitionShardPlan.java | 128 +++
.../sdk/io/iceberg/cdc/sink/SplitLateData.java | 102 ++
.../beam/sdk/io/iceberg/cdc/sink/TableSetup.java | 724 ++++++++++++++
.../beam/sdk/io/iceberg/CommitSchemaUnionTest.java | 780 +++++++++++++++
.../sdk/io/iceberg/cdc/sink/CommitTokenTest.java | 173 ++++
.../sdk/io/iceberg/cdc/sink/CommitWindowsTest.java | 411 ++++++++
.../iceberg/cdc/sink/PartitionShardPlanTest.java | 120 +++
.../sdk/io/iceberg/cdc/sink/TableSetupTest.java | 1028 ++++++++++++++++++++
34 files changed, 4876 insertions(+), 121 deletions(-)
copy
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/{package-info.java
=> FlushPayload.java} (55%)
create mode 100644
runners/spark/src/test/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingPipelineResultTest.java
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/CommitSchemaUnion.java
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CommitToken.java
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CommitWindows.java
copy
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/{ReadTaskDescriptor.java
=> cdc/sink/DestinationShard.java} (52%)
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/PartitionShardPlan.java
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/SplitLateData.java
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/TableSetup.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/CommitSchemaUnionTest.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CommitTokenTest.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CommitWindowsTest.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/PartitionShardPlanTest.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/TableSetupTest.java