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 14c61d13821 Bump zizmorcore/zizmor-action from 0.6.1 to 0.6.2 (#39604)
add 86ca05321d5 Adds the Delta Lake CDC read transforms to the Managed I/O
API (#39599)
add aa7f74cea25 mention otel in changes (#39618)
add 15082278fd1 Support sharded coder for Prism runner cross-lang (#39623)
add c2c910a70bf Fix Dataflow ValueProvider serialization (#39614)
add a0f3518d076 Deflake JmsIO tests (#39571)
add 16f471eee5c [Docs] Add a contributor guide for running Python on a
local Flink cluster (#39580)
add 7ec43af58eb fix PostCommit Python Dependency
add 7712ac62dd7 Merge pull request #39637 from
aIbrahiim/fix-postcommit-pydep-pyarrow
add 48c9e1f9b38 fix(dataframe): claim remaining restriction range on
empty/header-only CSV reads (#39581)
add c60b021cde4 [Iceberg CDC] Finish wiring CDC source together and add
external API (#39600)
add cb75d1b773f Updates CHANGES.md to include Delta Lake CDC
add 8a5d5d2f468 Merge pull request #39644 from
chamikaramj/update_change_log
add e793e8abab4 add Timestamp.MICROS for iceberg timestamptz (#39592)
add e2ae447bef7 Add google-api-python-client to Python 3.14 container
(#39640)
No new revisions were added by this update.
Summary of changes:
.../IO_Iceberg_Integration_Tests.json | 2 +-
.../IO_Iceberg_Integration_Tests_Dataflow.json | 2 +-
.../beam_PostCommit_Java_Delta_IO_Dataflow.json | 2 +-
.../beam_PostCommit_Python_Dependency.json | 4 +-
...m_PostCommit_Python_Xlang_Messaging_Direct.json | 2 +-
.github/trigger_files/beam_PostCommit_SQL.json | 2 +-
.github/trigger_files/beam_PreCommit_SQL.json | 2 +-
CHANGES.md | 7 +
contributor-docs/README.md | 1 +
contributor-docs/local-flink-python.md | 204 +++++
.../model/pipeline/v1/external_transforms.proto | 2 +
sdks/go/pkg/beam/runners/prism/internal/coders.go | 11 +
.../pkg/beam/runners/prism/internal/coders_test.go | 16 +
.../provider/iceberg/BeamSqlCliIcebergTest.java | 6 +-
.../meta/provider/iceberg/IcebergReadWriteIT.java | 3 -
.../sdk/extensions/sql/impl/rel/BeamCalcRel.java | 40 +
.../extensions/sql/impl/utils/CalciteUtils.java | 8 +-
.../sdk/extensions/sql/BeamComplexTypeTest.java | 32 +
.../extensions/sql/impl/rel/BeamCalcRelTest.java | 20 +
sdks/java/io/delta/build.gradle | 2 +-
.../beam/sdk/io/delta/DeltaCDCSourceDoFn.java | 8 +-
...va => DeltaCdcReadSchemaTransformProvider.java} | 78 +-
.../java/org/apache/beam/sdk/io/delta/DeltaIO.java | 48 +-
.../org/apache/beam/sdk/io/delta/DeltaIOIT.java | 168 +++-
.../org/apache/beam/sdk/io/delta/DeltaIOTest.java | 941 +++++++++++----------
.../beam/sdk/io/delta/DeltaWriteTestUtils.java | 371 ++++++++
.../IcebergCdcReadSchemaTransformProvider.java | 31 +-
.../org/apache/beam/sdk/io/iceberg/IcebergIO.java | 45 +-
.../beam/sdk/io/iceberg/IcebergScanConfig.java | 5 -
.../beam/sdk/io/iceberg/IncrementalScanSource.java | 102 ---
.../apache/beam/sdk/io/iceberg/ReadFromTasks.java | 98 ---
.../beam/sdk/io/iceberg/WatchForSnapshots.java | 190 -----
.../sdk/io/iceberg/cdc/ApplyWatermarkColumn.java | 99 +++
.../beam/sdk/io/iceberg/cdc/CdcOutputUtils.java | 29 +-
.../beam/sdk/io/iceberg/cdc/CdcResolver.java | 19 +-
.../io/iceberg/cdc/IncrementalChangelogSource.java | 211 +++++
.../beam/sdk/io/iceberg/cdc/LocalResolveDoFn.java | 8 +-
.../sdk/io/iceberg/cdc/ReadFromChangelogs.java | 11 +-
.../beam/sdk/io/iceberg/cdc/ResolveChanges.java | 168 ++++
.../beam/sdk/io/iceberg/cdc/SnapshotWindowFn.java | 87 ++
.../sdk/io/iceberg/cdc/WatchForSnapshotsSdf.java | 57 +-
.../IcebergCdcReadSchemaTransformProviderTest.java | 114 +++
.../beam/sdk/io/iceberg/IcebergScanConfigTest.java | 270 ++++++
.../IcebergSchemaTransformTranslationTest.java | 1 +
.../io/iceberg/catalog/IcebergCatalogBaseIT.java | 281 +++++-
.../io/iceberg/cdc/ApplyWatermarkColumnTest.java | 158 ++++
.../cdc/IncrementalChangelogSourceTest.java | 514 +++++++++++
.../sdk/io/iceberg/cdc/ResolveChangesTest.java | 222 +++++
.../sdk/io/iceberg/cdc/SnapshotWindowFnTest.java | 93 ++
.../io/iceberg/cdc/WatchForSnapshotsSdfTest.java | 49 +-
.../java/org/apache/beam/sdk/io/jms/JmsIOTest.java | 176 +---
.../org/apache/beam/sdk/io/jms/JmsLocalTest.java | 245 ++++++
.../java/org/apache/beam/sdk/managed/Managed.java | 5 +
sdks/python/apache_beam/dataframe/io.py | 1 +
sdks/python/apache_beam/dataframe/io_test.py | 6 +
.../runners/dataflow/internal/apiclient.py | 6 +-
.../runners/dataflow/internal/apiclient_test.py | 19 +
.../container/base_image_requirements_manual.txt | 2 +-
.../container/py314/base_image_requirements.txt | 2 +
sdks/python/test-suites/direct/build.gradle | 3 +-
sdks/python/test-suites/tox/py310/build.gradle | 27 +-
sdks/standard_expansion_services.yaml | 1 +
.../site/content/en/documentation/io/connectors.md | 10 +-
.../site/content/en/documentation/io/managed-io.md | 104 +++
64 files changed, 4341 insertions(+), 1110 deletions(-)
create mode 100644 contributor-docs/local-flink-python.md
copy
sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/{DeltaReadSchemaTransformProvider.java
=> DeltaCdcReadSchemaTransformProvider.java} (56%)
create mode 100644
sdks/java/io/delta/src/test/java/org/apache/beam/sdk/io/delta/DeltaWriteTestUtils.java
delete mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IncrementalScanSource.java
delete mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/ReadFromTasks.java
delete mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WatchForSnapshots.java
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/ApplyWatermarkColumn.java
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/IncrementalChangelogSource.java
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/ResolveChanges.java
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/SnapshotWindowFn.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergScanConfigTest.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/ApplyWatermarkColumnTest.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/IncrementalChangelogSourceTest.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/ResolveChangesTest.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/SnapshotWindowFnTest.java
create mode 100644
sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/JmsLocalTest.java