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 c280892881b [Iceberg CDC sink] Commit files stage (#40134)
add e4a60f588b0 [PubsubDynamicSink] fix hashmap order dependent test
(#40210)
add 5af8e2af167 Bump cloud.google.com/go/storage from 1.67.1 to 1.68.0 in
/sdks (#40218)
add a255f5a0d85 Bump github.com/nats-io/nats.go from 1.53.1 to 1.54.0 in
/sdks (#40217)
add fa134c16311 Bump github.com/aws/smithy-go from 1.28.1 to 1.28.2 in
/sdks (#40219)
add 7cfc6903c16 Bump github.com/nats-io/nats-server/v2 from 2.14.6 to
2.15.0 in /sdks (#40216)
add 7d0f1bf3870 Fix PickleCoder.as_deterministic_coder() raising TypeError
(#39943)
add cf2c5767e32 Kafka Streams runner: target flush markers across
repartition topics
add e038146a76b Kafka Streams runner: check flush multicast on a real
broker
add 25dc1ba1da8 Merge pull request #40186: Kafka Streams runner: target
flush markers across repartition topics
add 9863d92199e Bump org.nosphere.apache.rat from 0.9.0 to 0.11.0 (#40125)
add 28279d2e9f5 Extend GCS performance metrics and unify the GCS metric
namespace (#40142)
add 83148cba130 Portable tuple logical type (#40081)
add 07e8aaf226b fix Bigtable naming (#40206)
add 07d3169e4e0 Refactor load testing infra for upcoming GCS load tests
(#40223)
add c0f79b0251d [Spark][#36841] Add an opt-in Dataset-based backend for
portable pipelines (#40129)
add 1905e014069 Add GcsIOLoadTestBase and ParquetIO GCS load test (#40226)
add a40d8cd9bd9 Return last evaluated progress on lock timeout in
RestrictionTrackers.getProgress (#40205)
add 2906a89f93a [Iceberg CDC sink] Top API layer (#40161)
No new revisions were added by this update.
Summary of changes:
.../beam_PostCommit_Java_IO_Performance_Tests.json | 3 +-
...ommit_Java_PVR_Spark4_StructuredStreaming.json} | 3 +-
.github/workflows/README.md | 1 +
...Commit_Java_PVR_Spark4_StructuredStreaming.yml} | 20 +-
CHANGES.md | 1 +
build.gradle.kts | 6 +-
it/common/build.gradle | 1 +
.../common/dataflow/DefaultPipelineLauncher.java | 11 +-
.../apache/beam/it/common/utils/ByteSizeUtils.java | 114 ++
.../beam/it/common/dataflow/LoadTestBase.java | 50 +-
.../beam/it/common/storage/GcsIOLoadTestBase.java | 173 +++
.../beam/it/common/utils/ByteSizeUtilsTest.java | 89 ++
it/google-cloud-platform/build.gradle | 29 +-
.../apache/beam/it/gcp/storage/ParquetIOLT.java | 904 ++++++++++++
.../storage/{FileBasedIOLT.java => TextIOLT.java} | 10 +-
.../dataflow/worker/PubsubDynamicSinkTest.java | 58 +-
.../streams/translation/GroupByKeyTranslator.java | 14 +-
...tioner.java => KStreamsPayloadPartitioner.java} | 25 +-
.../streams/translation/ShuffleByKeyProcessor.java | 35 +-
.../streams/translation/TerminationTracker.java | 4 +-
.../KStreamsPayloadPartitionerBrokerIT.java | 176 +++
.../KStreamsPayloadPartitionerTest.java | 74 +
.../translation/ShuffleByKeyProcessorTest.java | 75 +-
runners/spark/4/README.md | 6 +
runners/spark/job-server/spark_job_server.gradle | 10 +-
.../beam/runners/spark/SparkPipelineOptions.java | 9 +
.../beam/runners/spark/SparkPipelineRunner.java | 18 +-
.../runners/spark/metrics/MetricsAccumulator.java | 12 +-
.../SparkDatasetPortablePipelineTranslator.java | 379 +++++
.../SparkDatasetTranslationContext.java | 126 ++
.../spark/SparkDatasetPortableExecutionTest.java} | 152 +-
...SparkDatasetPortablePipelineTranslatorTest.java | 442 ++++++
.../SparkDatasetTranslationContextTest.java | 164 +++
sdks/go.mod | 14 +-
sdks/go.sum | 28 +-
.../sdk/fn/splittabledofn/RestrictionTrackers.java | 66 +-
.../fn/splittabledofn/RestrictionTrackersTest.java | 55 +-
.../sdk/extensions/gcp/storage/GcsFileSystem.java | 157 +-
.../beam/sdk/extensions/gcp/util/GcsUtil.java | 14 +
.../beam/sdk/extensions/gcp/util/GcsUtilV1.java | 219 ++-
.../beam/sdk/extensions/gcp/util/Transport.java | 165 ++-
.../extensions/gcp/storage/GcsFileSystemTest.java | 61 +
.../beam/sdk/extensions/gcp/util/GcsUtilTest.java | 8 +-
.../sdk/extensions/gcp/util/TransportTest.java | 106 ++
.../java/org/apache/beam/sdk/io/text/TextIOIT.java | 8 +-
.../BigtableReadSchemaTransformProvider.java | 2 +-
.../BigtableWriteSchemaTransformProvider.java | 2 +-
.../org/apache/beam/sdk/io/iceberg/IcebergIO.java | 18 +
.../beam/sdk/io/iceberg/IcebergWriteResult.java | 87 ++
.../beam/sdk/io/iceberg/cdc/sink/WriteCdcRows.java | 550 +++++++
.../beam/sdk/io/iceberg/cdc/sink/package-info.java | 233 ++-
.../iceberg/cdc/sink/WriteCdcRowsReadBackTest.java | 1526 ++++++++++++++++++++
.../sdk/io/iceberg/cdc/sink/WriteCdcRowsTest.java | 1499 +++++++++++++++++++
sdks/python/apache_beam/coders/coder_impl.pxd | 1 +
sdks/python/apache_beam/coders/coder_impl.py | 86 +-
sdks/python/apache_beam/coders/coders.py | 33 +-
sdks/python/apache_beam/coders/coders_test.py | 15 +
sdks/python/apache_beam/coders/row_coder.py | 8 +
sdks/python/apache_beam/coders/row_coder_test.py | 49 +
.../src/yaml/EmojiMap.ts | 3 +-
sdks/python/apache_beam/transforms/sql_test.py | 34 +-
.../typehints/native_type_compatibility.py | 20 +-
.../typehints/native_type_compatibility_test.py | 14 +
sdks/python/apache_beam/typehints/row_type_test.py | 12 +
sdks/python/apache_beam/typehints/schemas.py | 209 ++-
sdks/python/apache_beam/typehints/schemas_test.py | 11 +-
sdks/python/apache_beam/yaml/standard_io.yaml | 15 +-
sdks/python/apache_beam/yaml/tests/bigtable.yaml | 8 +-
68 files changed, 8095 insertions(+), 435 deletions(-)
copy .github/trigger_files/{beam_PostCommit_Python_Versions.json =>
beam_PostCommit_Java_PVR_Spark4_StructuredStreaming.json} (66%)
copy .github/workflows/{beam_PostCommit_Java_PVR_Spark4_Streaming.yml =>
beam_PostCommit_Java_PVR_Spark4_StructuredStreaming.yml} (89%)
create mode 100644
it/common/src/main/java/org/apache/beam/it/common/utils/ByteSizeUtils.java
create mode 100644
it/common/src/test/java/org/apache/beam/it/common/storage/GcsIOLoadTestBase.java
create mode 100644
it/common/src/test/java/org/apache/beam/it/common/utils/ByteSizeUtilsTest.java
create mode 100644
it/google-cloud-platform/src/test/java/org/apache/beam/it/gcp/storage/ParquetIOLT.java
rename
it/google-cloud-platform/src/test/java/org/apache/beam/it/gcp/storage/{FileBasedIOLT.java
=> TextIOLT.java} (96%)
rename
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/{GroupByKeyBroadcastPartitioner.java
=> KStreamsPayloadPartitioner.java} (74%)
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/KStreamsPayloadPartitionerBrokerIT.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/KStreamsPayloadPartitionerTest.java
create mode 100644
runners/spark/src/main/java/org/apache/beam/runners/spark/translation/SparkDatasetPortablePipelineTranslator.java
create mode 100644
runners/spark/src/main/java/org/apache/beam/runners/spark/translation/SparkDatasetTranslationContext.java
copy
runners/{flink/src/test/java/org/apache/beam/runners/flink/PortableExecutionTest.java
=>
spark/src/test/java/org/apache/beam/runners/spark/SparkDatasetPortableExecutionTest.java}
(50%)
create mode 100644
runners/spark/src/test/java/org/apache/beam/runners/spark/translation/SparkDatasetPortablePipelineTranslatorTest.java
create mode 100644
runners/spark/src/test/java/org/apache/beam/runners/spark/translation/SparkDatasetTranslationContextTest.java
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/WriteCdcRows.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/WriteCdcRowsReadBackTest.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/WriteCdcRowsTest.java