This is an automated email from the ASF dual-hosted git repository.
derrickaw pushed a change to branch 20260722_removeGsutil
in repository https://gitbox.apache.org/repos/asf/beam.git
omit 6351fa15658 fix influxdb errors
omit 91f90c1aaf2 add workflow to table
omit a5a5f2fa6da add comment
omit 4fced04bee3 remove gsutil usage
add abe3fc77c3d Fix silent per-char iteration when DoFn returns
str/bytes/dict (#38429)
add a15b880c53c Merge pull request #39043: Improve WithKeys coder
inference context
add b04ca1e0cee Fix pyrefly check bad-specialization (#39419)
add c8a689f7c0d Fix pyrefly check bad-argument-count (#39416)
add da7693c3dd1 [Gemini] Fix pyrefly check invalid-yield (#39417)
add f1d9c3e122e docs: Document pipeline runner API protos (#39292)
add adbb11ddb46 Fix gcs endpoint pipeline option wiring (#39435)
add 5e499388670 KeyCommitTooLarge logging improvements (#39316)
add 7356d4be850 Replace port with listeners in LocalKafka configuration
add 4f0044b5de9 Merge pull request #39438 from
sjvanrossum/testing-kafka-service-fix
add c7ab90ebad3 Bump google.golang.org/api from 0.289.0 to 0.290.0 in
/sdks (#39449)
add b0a7099e1f1 Fix lz4 conflicts (#39454)
add 1c9fd1bc1b9 Remove isort dependency, migrate to ruff equivalent
(#39411)
add c8aa28f7e2e docs: Document Fn Execution API protos (#39293)
add d688e1b9505 [Java SDK] Warn when ValueState contains collection types
(#37530)
add a72d6cba37a [KafkaIO] Use consumer position and lag to estimate end
offsets (#39285)
add 0140552d6d5 [#30019] Fix type checking failures around unions with
nested complex types (#39413)
add 9833974946c Remove legacy HttpError usage in Bigtable I/O IT test
(#39441)
add 9c3f51e56e0 [Python] Add Watch transform with growth_of polling SDF
(#39023)
add eb9e7a3dacd Override default fadvise to fix regression from
gcs-connector v3 upgrade (#39445)
add a3703d38100 [Dataflow Streaming][Multikey] Support MultiKey commits in
windmill clients (#38768)
add b15e5f47d40 Add equal_to_approx matcher for approximate numeric
assertions (#39443)
add 7f61891553c Bump cloud.google.com/go/datastore from 1.25.0 to 1.26.0
in /sdks (#39468)
add 3cf04be0dff Bump cloud.google.com/go/bigtable from 1.50.0 to 1.51.0 in
/sdks (#39469)
add 1b1b70645fb Bump docker/login-action from 4.4.0 to 4.5.0 (#39470)
add 2232808083b Bump zizmorcore/zizmor-action from 0.6.0 to 0.6.1 (#39471)
add 60455c3568a Bump setuptools from 78.1.1 to 83.0.0 in
/.test-infra/mock-apis (#39482)
add 0ffb108ab79 Revert "Merge pull request #39043: Improve WithKeys coder
inference context" (#39483)
add 72d7ae3b414 update tour of beam workflow go version (#39486)
add d6ab9929a4e Bump docker/login-action from 4.5.0 to 4.5.1 (#39497)
add cb7329c3273 Use prebuilt Snapshots SDK images for PostCommit Python
Arm (#39498)
add 6741628d3af Pin Playground kafka-emulator to kafka-clients 2.4.1
(#39503)
add 226ad87e0d9 enable otel context propagation - runner v1 sink, source
changes, doFnRunner changes for per element propagation (#39152)
add 6d752f569d9 OTEL in spanner. (#39149)
add 70a522364e9 (IcebergIO) Support PartitionSpec/SortOrder on dynamic
table creation via IcebergIO (#39408)
add faa3ad81a95 Skip IcebergPerformanceTest until next release
add d9d897a37f1 Merge pull request #39502 from apache/skip-iceberg-perf
add 214863b2e04 Bump golang.org/x/net from 0.54.0 to 0.55.0 in
/.test-infra/mock-apis (#39510)
add de234c72b88 Fix inconsistent AvroSchema type and value for
SqlType.Date values (#39414)
add ac5262370d9 Adds documentation for the Delta Lake Read Managed I/O
(#39495)
add 901fccd7cf6 sdks/java: remove DefaultAnnotation(NonNull) from
package-info.java files
add 1d25d3b75a4 Merge pull request #39463: Remove
@DefaultAnnotation(NonNull.class) throughout project - it is already default
add fdac189a032 Fix nullness for PubsubIO
add 405438b83fe Merge pull request #39444: Fix nullness for PubsubIO
add dec8d23717a Bump torch (#39512)
add 41609956aa3 OTEL in kafka. (#39151)
add ec93d37c602 Bump golang.org/x/oauth2 from 0.7.0 to 0.27.0 in
/playground/backend (#39524)
add ad1e278da3e sdks/java: re-enable nullness checks in WithKeys (#39506)
add 9ae08d5a068 [Solace] Close the HTTP response content stream in
BrokerResponse (#39404)
add 9c82e053f02 Avoid output inside try-catch in Java IO (#39124)
add dd463fa934d Fix flaky AsyncWrapper reset_state test on Python 3.14
(#39521)
add 2eb3323c555 Bump github.com/moby/moby/client from 0.5.0 to 0.5.1 in
/sdks (#39517)
add 905ade4ed42 Bump github.com/aws/smithy-go from 1.27.4 to 1.27.5 in
/sdks (#39519)
add cae3e1749f2 Bump actions/stale from 10 to 11 (#39520)
add f456c02459d Bump scikit-learn (#39525)
add 39c0dde7080 [Gemini] Fix pyrefly check bad-typed-dict-key (#39415)
add 24346193cc8 Add registerSqlOperator() to BeamSqlEnv for custom SQL
operators (#39432)
add f258e3e8ae4 OTEL in pubsub (#39150)
add 55fde074af8 Add the directory with staged files to sys.path and
document the usage (#39434)
add fc0d9895080 Replace non-PEP 585 types in watch.py (#39527)
add f20da8a1c4a Update ruff and pyrefly dependencies (#39531)
add 8473d90e93f [Java IO] Add ArrowFlight IO connector (#37904)
add 2ee432b2b61 Support JmsIO SchemaTransform and cross-lang (#39437)
add 3926b886590 fix golangci-lint issue - tour of beam (#39490)
add 4fd1744935c [Gemini] Fix pyrefly check unexpected-keyword (#39528)
add 509af44c580 Bump docker/login-action from 4.5.1 to 4.5.2 (#39540)
add e9d9086ff01 Bump github.com/aws/aws-sdk-go-v2 from 1.43.0 to 1.43.1 in
/sdks (#39542)
add c4a6f802b74 Bump github.com/aws/aws-sdk-go-v2/service/s3 in /sdks
(#39541)
add 79504e2e6a2 Bump google.golang.org/api from 0.290.0 to 0.291.0 in
/sdks (#39544)
add 4738b16cf40 Bump github.com/aws/aws-sdk-go-v2/credentials in /sdks
(#39543)
add b5c6009c057 Bump github.com/aws/aws-sdk-go-v2/feature/s3/manager in
/sdks (#39545)
add f6c7a06d90a Fix Update Python Dependencies (#39534)
add 3c9f9ce84df [Dataflow Streaming] [Multi Key] MultiKey failure handling
+ Integration (#38919)
add 9019efd790c add workflow_dispatch to other workflow files (#39491)
add 08dc50a3b91 Fix flaky BigQuery persistent retry test (#39539)
add 8a1c19bc9d0 [Python] Convert typing and native generic hints in Watch
coder inference (#39547)
add 6e2044b0083 Bump lower and upper bounds for pyarrow + related
dependencies, remove unnecessary CVE hotfix (#39530)
add eda08d8a3c6 Changes SplittableDoFn to call TruncateRestriction on
drain (#39535)
add ec7004f6f57 [Interactive Beam] Fix caching deadlock, wait race
conditions, and stale graph in notebooks (#39161)
add 58bac320ebd [IcebergIO] Raise Java 17 floor for IcebergIO's Java 11
dependents (#39064)
add 2c735b67eb3 Bump cloud.google.com/go/spanner from 1.93.0 to 1.94.0 in
/sdks (#39550)
add c6e0f5b630e Bump github.com/aws/aws-sdk-go-v2/config in /sdks (#39551)
add 2fd6af992b5 Bump github.com/aws/aws-sdk-go-v2/service/s3 in /sdks
(#39554)
add 61fad4f672f Bump github.com/aws/aws-sdk-go-v2/feature/s3/manager in
/sdks (#39553)
add 525780e7af7 Provide a better error when beam plugin was supplied but
wasn't staged. (#39440)
add f65e0e0b0a8 Updates the Delta Lake source to support reading bounded
change data (#39426)
add 4f934835059 Fix module-level side effects and global random seeding in
univariate ML anomaly tests (#39462)
add bc991d99d17 Set envs
add 8a7c972c61b Merge pull request #39558 from apache/fix-auditkeys
add 621a78cd412 Add Beam YAML support for DebeziumIO
add 01c668b1455 Add Debezium YAML integration test
add a2267554a16 Fix record schema
add 25ed2f8b5b7 With max num of records
add aae48ec3ff0 Add setters
add 2541fc77cb3 Refactoring
add 315ed95a97d Fix python formatter
add b85646bf8df Remove primaryKeyColumns options
add 5178343b652 Fix spotless
add 0de9a676681 Merge pull request #39457 from apache/debezium-io-yaml
add b9e4e2f27a8 Bump docker/login-action from 4.5.2 to 4.6.0 (#39552)
add f7d8d7c8b58 Preserve partitioning on temp FILE_LOADS tables (#38833)
add 141804ab568 fix AddFilesIT filter for BigLake (#39533)
add f5feab8c598 Enhance Python Timestamp to be precision-variable up to
nanos, and map it to Timestamp logical type (#39537)
add 7b9380b1c25 Buffer BufferedLogger by newline to avoid log splitting
(#39288)
add 03db1a07096 Fix DataflowOutputCounter calculation for
ValueInEmptyWindows (#39487)
add b8d77b86055 [IcebergIO] Upgrade Iceberg dependency to 1.11.0 (#39559)
add 980c11432a1 Clean up legacy references to apitools in GCS I/O (#39433)
add 9c561e2983e [Iceberg] Make timestamptz return new Timestamp.MICROS
logical type (#39344)
add 459b7ee5036 Fix flaky unit test to pass post-submit checks (#39562)
add 42c693f3c1e Bump github/codeql-action from 4 to 4.37.3 (#39564)
add 5c58b58cff6 Support IBM MQ for Python JmsIO (#39467)
add 0619156f8b5 use Java 17 harness (#39570)
add 2621e9e047c Remove remaining artifacts from dataflow apitools client
(#39439)
add 51966762894 update containers (#39575)
add e4779cf79f2 Fix flaky FileIOTest.testMatchWatchForNewFiles test under
CI filesystems (#38047)
add 7f96ee4f2bf Bump google.golang.org/grpc from 1.82.1 to 1.83.0 in /sdks
(#39585)
add 93af6786bbb Bump github/codeql-action from 4.37.3 to 4.37.4 (#39586)
add 8b4c6751963 Support core dump analysis with pystack and gdb. (#39484)
add a72451d8072 Bump github.com/nats-io/nats-server/v2 from 2.14.3 to
2.14.4 in /sdks (#39584)
add 3a02af8146c [Docs] Update Flink version references on the Flink runner
page (#39212)
add 539b048eb8d Fix dataframe CSV tests on Windows (#39563)
add 89a3d5d0dc1 remove gsutil usage
add 0fb49e6ffc8 add comment
add 51c9b1d8923 add workflow to table
add 7560d432083 fix influxdb errors
This update added new revisions after undoing existing revisions.
That is to say, some revisions that were in the old version of the
branch are not in the new version. This situation occurs
when a user --force pushes a change and generates a repository
containing something like this:
* -- * -- B -- O -- O -- O (6351fa15658)
\
N -- N -- N refs/heads/20260722_removeGsutil (7560d432083)
You should already have received notification emails for all of the O
revisions, and so the following emails describe only the N revisions
from the common base, B.
Any revisions marked "omit" are not gone; other references still
refer to them. Any revisions marked "discard" are gone forever.
No new revisions were added by this update.
Summary of changes:
.../actions/setup-environment-action/action.yml | 24 +-
.../IO_Iceberg_Integration_Tests.json | 2 +-
.github/trigger_files/beam_PostCommit_Python.json | 2 +-
..._Spark.json => beam_PostCommit_Python_Arm.json} | 0
.../beam_PostCommit_Python_Xlang_Gcp_Direct.json | 2 +-
.../beam_PostCommit_Python_Xlang_IO_Dataflow.json | 2 +-
.../beam_PostCommit_Python_Xlang_IO_Direct.json | 2 +-
...m_PostCommit_Python_Xlang_Messaging_Direct.json | 2 +-
.../beam_Infrastructure_AuditUnmanagedKeys.yml | 5 +
.../beam_PerformanceTests_xlang_KafkaIO_Python.yml | 2 +-
.../beam_PostCommit_Java_IO_Performance_Tests.yml | 2 +-
.github/workflows/beam_PostCommit_Python_Arm.yml | 21 +-
.../beam_PostCommit_Python_Xlang_IO_Dataflow.yml | 2 +-
.../beam_PostCommit_Python_Xlang_IO_Direct.yml | 2 +-
.github/workflows/beam_PreCommit_GHA.yml | 2 +-
.github/workflows/build_release_candidate.yml | 2 +-
.github/workflows/codeql.yml | 4 +-
.github/workflows/finalize_release.yml | 2 +-
.github/workflows/python_dependency_tests.yml | 1 +
.github/workflows/stale.yml | 2 +-
.github/workflows/tour_of_beam_backend.yml | 10 +-
.../workflows/tour_of_beam_backend_integration.yml | 1 +
.github/workflows/update_python_dependencies.yml | 2 +
.test-infra/mock-apis/go.mod | 2 +-
.test-infra/mock-apis/go.sum | 4 +-
.test-infra/mock-apis/poetry.lock | 18 +-
CHANGES.md | 18 +-
build.gradle.kts | 1 +
.../org/apache/beam/gradle/BeamModulePlugin.groovy | 13 +-
examples/java/iceberg/build.gradle | 4 +-
it/build.gradle | 2 +-
it/iceberg/build.gradle | 4 +-
learning/tour-of-beam/backend/function.go | 6 +-
.../backend/integration_tests/client.go | 8 +-
.../backend/internal/fs_content/yaml.go | 10 +-
.../backend/internal/storage/datastore.go | 11 +-
.../tour-of-beam/backend/internal/storage/mock.go | 8 +-
.../beam/model/fn_execution/v1/beam_fn_api.proto | 63 +-
.../model/fn_execution/v1/beam_provision_api.proto | 2 +-
.../beam/model/fnexecution/v1/standard_coders.yaml | 31 +
.../beam/model/pipeline/v1/beam_runner_api.proto | 94 ++-
.../apache/beam/model/pipeline/v1/endpoints.proto | 3 +
.../model/pipeline/v1/external_transforms.proto | 20 +
.../apache/beam/model/pipeline/v1/metrics.proto | 73 +-
.../org/apache/beam/model/pipeline/v1/schema.proto | 107 +++
playground/backend/go.mod | 7 +-
playground/backend/go.sum | 10 +-
playground/kafka-emulator/build.gradle | 11 +
runners/core-java/build.gradle | 1 +
.../core/GroupAlsoByWindowViaWindowSetNewDoFn.java | 2 +-
.../apache/beam/runners/core/KeyedWorkItem.java | 9 +
.../apache/beam/runners/core/ReduceFnRunner.java | 18 +-
.../apache/beam/runners/core/SimpleDoFnRunner.java | 13 +
.../core/SplittableParDoViaKeyedWorkItems.java | 70 +-
.../runners/core/construction/package-info.java | 4 -
.../beam/runners/core/metrics/package-info.java | 4 -
.../org/apache/beam/runners/core/package-info.java | 4 -
.../beam/runners/core/triggers/package-info.java | 4 -
.../runners/core/SplittableParDoProcessFnTest.java | 283 ++++++-
.../runners/extensions/metrics/package-info.java | 4 -
runners/google-cloud-dataflow-java/build.gradle | 4 +-
.../google-cloud-dataflow-java/worker/build.gradle | 1 +
.../dataflow/worker/DataflowExecutionContext.java | 4 +
.../dataflow/worker/DataflowOutputCounter.java | 68 +-
.../worker/IntrinsicMapTaskExecutorFactory.java | 11 +-
.../dataflow/worker/MultiKeyBundleOptions.java | 138 ++++
.../dataflow/worker/SimpleParDoFnHelpers.java | 5 +-
.../dataflow/worker/StreamingDataflowWorker.java | 36 +-
.../StreamingGroupAlsoByWindowViaWindowSetFn.java | 2 +-
.../worker/StreamingModeExecutionContext.java | 127 +++-
.../dataflow/worker/UngroupedWindmillReader.java | 7 +-
.../dataflow/worker/WindmillKeyedWorkItem.java | 31 +-
.../WindmillOpenTelemetryContextPropagator.java | 22 +-
.../beam/runners/dataflow/worker/WindmillSink.java | 29 +-
...Exception.java => WorkCancellingException.java} | 30 +-
.../worker/WorkItemCancelledException.java | 29 +-
.../dataflow/worker/streaming/ActiveWorkState.java | 5 +
.../streaming/BoundedQueueExecutorWorkHandle.java | 7 +-
.../worker/streaming/ComputationState.java | 7 +
.../worker/streaming/ComputationWorkExecutor.java | 5 +-
.../streaming/KeyCommitTooLargeException.java | 20 +-
.../runners/dataflow/worker/streaming/Work.java | 60 +-
.../dataflow/worker/util/BoundedQueueExecutor.java | 36 +-
.../worker/windmill/client/WindmillStream.java | 5 +
.../worker/windmill/client/commits/Commit.java | 72 +-
.../worker/windmill/client/commits/Commits.java | 2 +-
.../windmill/client/commits/CompleteCommit.java | 30 +-
.../commits/StreamingApplianceWorkCommitter.java | 13 +-
.../commits/StreamingEngineWorkCommitter.java | 113 ++-
.../client/getdata/StreamGetDataClient.java | 5 +-
.../windmill/client/grpc/GrpcCommitWorkStream.java | 93 ++-
.../worker/windmill/state/WindmillStateReader.java | 9 +-
.../processing/ComputationWorkExecutorFactory.java | 23 +-
.../work/processing/StreamingWorkScheduler.java | 234 +++---
.../processing/failures/WorkFailureProcessor.java | 73 +-
.../dataflow/worker/DataflowOutputCounterTest.java | 108 +++
.../dataflow/worker/FakeWindmillServer.java | 67 +-
.../IntrinsicMapTaskExecutorFactoryTest.java | 14 +-
.../worker/KeyTokenInvalidExceptionTest.java | 39 -
.../worker/StreamingDataflowWorkerTest.java | 383 ++++++++--
.../worker/StreamingModeExecutionContextTest.java | 291 ++++++-
.../worker/WindmillReaderIteratorBaseTest.java | 4 +-
.../worker/WindowingWindmillReaderTest.java | 3 +-
.../dataflow/worker/WorkerCustomSourcesTest.java | 18 +-
.../worker/streaming/ActiveWorkStateTest.java | 33 +-
.../streaming/ComputationStateCacheTest.java | 4 +-
.../worker/streaming/ComputationStateTest.java | 114 +++
.../dataflow/worker/streaming/WorkTest.java | 4 +-
.../worker/util/BoundedQueueExecutorTest.java | 82 +-
.../worker/util/KeyGroupWorkQueueTest.java | 10 +-
.../StreamingApplianceWorkCommitterTest.java | 15 +-
.../commits/StreamingEngineWorkCommitterTest.java | 323 +++++++-
.../client/grpc/GrpcCommitWorkStreamTest.java | 273 +++++++
.../windmill/state/WindmillStateReaderTest.java | 10 +-
.../failures/WorkFailureProcessorTest.java | 101 ++-
.../work/refresh/ActiveWorkRefresherTest.java | 4 +-
.../worker/windmill/src/main/proto/windmill.proto | 23 +
sdks/go.mod | 68 +-
sdks/go.sum | 136 ++--
sdks/go/container/tools/buffered_logging.go | 64 +-
sdks/go/container/tools/buffered_logging_test.go | 168 +++-
.../prism/internal/engine/elementmanager.go | 1 +
.../org/apache/beam/sdk/jmh/util/package-info.java | 4 -
.../apache/beam/sdk/annotations/package-info.java | 4 -
.../org/apache/beam/sdk/coders/package-info.java | 4 -
.../apache/beam/sdk/expansion/package-info.java | 4 -
.../beam/sdk/fn/splittabledofn/package-info.java | 4 -
.../org/apache/beam/sdk/harness/package-info.java | 3 -
.../org/apache/beam/sdk/io/fs/package-info.java | 4 -
.../java/org/apache/beam/sdk/io/package-info.java | 4 -
.../org/apache/beam/sdk/io/range/package-info.java | 4 -
.../org/apache/beam/sdk/metrics/package-info.java | 4 -
.../java/org/apache/beam/sdk/package-info.java | 4 -
.../org/apache/beam/sdk/runners/package-info.java | 3 -
.../apache/beam/sdk/schemas/SchemaTranslation.java | 2 +
.../beam/sdk/schemas/annotations/package-info.java | 4 -
.../apache/beam/sdk/schemas/io/package-info.java | 4 -
.../beam/sdk/schemas/io/payloads/package-info.java | 4 -
.../beam/sdk/schemas/logicaltypes/Timestamp.java | 9 +-
.../sdk/schemas/logicaltypes/package-info.java | 4 -
.../org/apache/beam/sdk/schemas/package-info.java | 4 -
.../sdk/schemas/parser/generated/package-info.java | 4 -
.../beam/sdk/schemas/parser/package-info.java | 4 -
.../beam/sdk/schemas/transforms/package-info.java | 4 -
.../schemas/transforms/providers/package-info.java | 4 -
.../beam/sdk/schemas/utils/package-info.java | 4 -
.../org/apache/beam/sdk/state/package-info.java | 4 -
.../org/apache/beam/sdk/testing/package-info.java | 4 -
.../org/apache/beam/sdk/transforms/WithKeys.java | 27 +-
.../beam/sdk/transforms/display/package-info.java | 4 -
.../sdk/transforms/errorhandling/package-info.java | 4 -
.../beam/sdk/transforms/join/package-info.java | 4 -
.../apache/beam/sdk/transforms/package-info.java | 4 -
.../sdk/transforms/reflect/DoFnSignatures.java | 49 ++
.../beam/sdk/transforms/reflect/package-info.java | 3 -
.../transforms/splittabledofn/package-info.java | 4 -
.../sdk/transforms/windowing/package-info.java | 4 -
.../sdk/util/construction/graph/package-info.java | 4 -
.../beam/sdk/util/construction/package-info.java | 4 -
.../sdk/values/OpenTelemetryContextPropagator.java | 8 +-
.../org/apache/beam/sdk/values/WindowedValues.java | 8 +-
.../org/apache/beam/sdk/values/package-info.java | 4 -
.../java/org/apache/beam/sdk/io/FileIOTest.java | 22 +-
.../beam/sdk/schemas/SchemaTranslationTest.java | 7 +
.../sdk/transforms/reflect/DoFnSignaturesTest.java | 117 +++
.../beam/sdk/extensions/arrow/ArrowConversion.java | 61 ++
.../sdk/extensions/arrow/ArrowConversionTest.java | 33 +
.../sdk/extensions/avro/coders/package-info.java | 4 -
.../beam/sdk/extensions/avro/io/package-info.java | 4 -
.../beam/sdk/extensions/avro/package-info.java | 4 -
.../avro/schemas/io/payloads/package-info.java | 4 -
.../sdk/extensions/avro/schemas/package-info.java | 4 -
.../extensions/avro/schemas/utils/AvroUtils.java | 2 +-
.../avro/schemas/utils/package-info.java | 4 -
.../avro/schemas/utils/AvroUtilsTest.java | 4 +-
.../google-cloud-platform-core/build.gradle | 1 +
.../sdk/extensions/gcp/options/GcsOptions.java | 85 ++-
.../beam/sdk/extensions/gcp/util/GcsUtilV1.java | 17 +-
.../sdk/extensions/gcp/GcpCoreApiSurfaceTest.java | 2 +
.../sdk/extensions/gcp/options/GcsOptionsTest.java | 47 ++
.../beam/sdk/extensions/gcp/util/GcsUtilTest.java | 13 +
sdks/java/extensions/sql/iceberg/build.gradle | 4 +-
.../beam/sdk/extensions/sql/impl/BeamSqlEnv.java | 13 +
.../extensions/sql/impl/CalciteQueryPlanner.java | 10 +-
.../sdk/extensions/sql/impl/JdbcConnection.java | 22 +
.../sdk/extensions/sql/impl/rel/package-info.java | 4 -
.../sdk/extensions/sql/impl/rule/package-info.java | 4 -
.../sql/impl/transform/agg/package-info.java | 4 -
.../sql/meta/provider/mongodb/package-info.java | 4 -
.../sql/meta/provider/pubsub/package-info.java | 4 -
.../sql/impl/BeamSqlEnvRegisterOperatorTest.java | 100 +++
.../io/{synthetic => arrow-flight}/build.gradle | 28 +-
.../beam/sdk/io/arrowflight/ArrowFlightIO.java | 840 ++++++++++++++++++++
.../beam/sdk/io/arrowflight}/package-info.java | 13 +-
.../beam/sdk/io/arrowflight/ArrowFlightIOTest.java | 330 ++++++++
.../org/apache/beam/io/debezium/DebeziumIO.java | 16 +-
.../DebeziumReadSchemaTransformProvider.java | 96 ++-
.../DebeziumReadSchemaTransformProviderTest.java | 139 ++++
.../beam/sdk/io/delta/CreateCDCReadTasksDoFn.java | 296 ++++++++
.../apache/beam/sdk/io/delta/DeltaCDCReadTask.java | 125 +++
.../beam/sdk/io/delta/DeltaCDCSourceDoFn.java | 359 +++++++++
.../java/org/apache/beam/sdk/io/delta/DeltaIO.java | 157 +++-
.../io/delta/DeltaReadSchemaTransformProvider.java | 6 +-
.../apache/beam/sdk/io/delta/DeltaSourceDoFn.java | 2 +-
.../org/apache/beam/sdk/io/delta/DeltaIOTest.java | 841 +++++++++++++++++++--
sdks/java/io/google-cloud-platform/build.gradle | 2 +
.../beam/sdk/io/gcp/bigquery/BigQueryIO.java | 17 +-
.../beam/sdk/io/gcp/bigquery/BigQueryUtils.java | 6 +
.../apache/beam/sdk/io/gcp/healthcare/FhirIO.java | 6 +-
.../apache/beam/sdk/io/gcp/healthcare/HL7v2IO.java | 17 +-
.../sdk/io/gcp/pubsub/AddTimestampAttribute.java | 13 +-
.../beam/sdk/io/gcp/pubsub/ExternalWrite.java | 15 +-
.../beam/sdk/io/gcp/pubsub/NestedRowToMessage.java | 8 +-
.../io/gcp/pubsub/PubSubPayloadTranslation.java | 46 +-
.../beam/sdk/io/gcp/pubsub/PubsubClient.java | 59 +-
.../beam/sdk/io/gcp/pubsub/PubsubGrpcClient.java | 47 +-
.../apache/beam/sdk/io/gcp/pubsub/PubsubIO.java | 301 ++++++--
.../beam/sdk/io/gcp/pubsub/PubsubJsonClient.java | 100 ++-
.../beam/sdk/io/gcp/pubsub/PubsubMessage.java | 23 +-
.../beam/sdk/io/gcp/pubsub/PubsubMessageToRow.java | 42 +-
...hAttributesAndMessageIdAndOrderingKeyCoder.java | 19 +-
...bsubMessageWithAttributesAndMessageIdCoder.java | 14 +-
.../pubsub/PubsubMessageWithAttributesCoder.java | 9 +-
.../pubsub/PubsubMessageWithMessageIdCoder.java | 9 +-
.../pubsub/PubsubReadSchemaTransformProvider.java | 34 +-
.../beam/sdk/io/gcp/pubsub/PubsubRowToMessage.java | 58 +-
.../sdk/io/gcp/pubsub/PubsubSchemaIOProvider.java | 51 +-
.../beam/sdk/io/gcp/pubsub/PubsubTestClient.java | 118 +--
.../sdk/io/gcp/pubsub/PubsubUnboundedSink.java | 62 +-
.../sdk/io/gcp/pubsub/PubsubUnboundedSource.java | 190 +++--
.../pubsub/PubsubWriteSchemaTransformProvider.java | 59 +-
.../apache/beam/sdk/io/gcp/pubsub/TestPubsub.java | 92 ++-
.../beam/sdk/io/gcp/pubsub/TestPubsubSignal.java | 79 +-
.../beam/sdk/io/gcp/spanner/BatchSpannerRead.java | 13 +-
.../sdk/io/gcp/spanner/CreateTransactionFn.java | 8 +-
.../beam/sdk/io/gcp/spanner/NaiveSpannerRead.java | 8 +-
.../beam/sdk/io/gcp/spanner/ReadSpannerSchema.java | 8 +-
.../beam/sdk/io/gcp/spanner/SpannerAccessor.java | 32 +-
.../beam/sdk/io/gcp/spanner/SpannerConfig.java | 15 +
.../apache/beam/sdk/io/gcp/spanner/SpannerIO.java | 42 +-
.../gcp/spanner/changestreams/dao/DaoFactory.java | 16 +-
.../dofn/CleanUpReadChangeStreamDoFn.java | 7 +
.../dofn/DetectNewPartitionsDoFn.java | 5 +-
.../spanner/changestreams/dofn/InitializeDoFn.java | 7 +
.../dofn/ReadChangeStreamPartitionDoFn.java | 5 +-
.../apache/beam/sdk/io/gcp/GcpApiSurfaceTest.java | 1 +
.../sdk/io/gcp/bigquery/BigQueryUtilsTest.java | 42 +-
.../sdk/io/gcp/spanner/SpannerIOWriteTest.java | 14 +-
.../dofn/ReadChangeStreamPartitionDoFnTest.java | 3 +-
sdks/java/io/hadoop-format/build.gradle | 1 +
sdks/java/io/iceberg/build.gradle | 10 +-
.../beam/sdk/io/iceberg/DynamicDestinations.java | 12 +-
.../org/apache/beam/sdk/io/iceberg/IcebergIO.java | 47 +-
.../beam/sdk/io/iceberg/IcebergScanConfig.java | 9 +-
.../apache/beam/sdk/io/iceberg/IcebergUtils.java | 76 +-
.../beam/sdk/io/iceberg/IncrementalScanSource.java | 4 +-
.../io/iceberg/OneTableDynamicDestinations.java | 33 +-
.../apache/beam/sdk/io/iceberg/ReadFromTasks.java | 4 +-
.../org/apache/beam/sdk/io/iceberg/ScanSource.java | 4 +-
.../apache/beam/sdk/io/iceberg/ScanTaskReader.java | 5 +-
.../io/iceberg/WritePartitionedRowsToFiles.java | 13 +-
.../org/apache/beam/sdk/io/iceberg/AddFilesIT.java | 15 +-
.../beam/sdk/io/iceberg/IcebergIOReadTest.java | 48 ++
.../beam/sdk/io/iceberg/IcebergIOWriteTest.java | 54 ++
.../beam/sdk/io/iceberg/IcebergUtilsTest.java | 43 +-
.../IcebergWriteSchemaTransformProviderTest.java | 13 +-
.../catalog/BigQueryMetastoreCatalogIT.java | 1 +
.../io/iceberg/catalog/IcebergCatalogBaseIT.java | 11 +-
sdks/java/io/jms/build.gradle | 6 +
...r.java => BeamGenericJmsConnectionFactory.java} | 28 +-
.../beam/sdk/io/jms/ConnectionConfiguration.java | 252 ++++++
.../java/org/apache/beam/sdk/io/jms/JmsIO.java | 30 +
.../io/jms/JmsReadSchemaTransformProvider.java} | 125 ++-
.../io/jms/JmsWriteSchemaTransformProvider.java | 191 +++++
.../sdk/io/jms/ConnectionConfigurationTest.java | 159 ++++
.../sdk/io/jms/JmsSchemaTransformProviderTest.java | 253 +++++++
sdks/java/io/kafka/build.gradle | 2 +
.../java/org/apache/beam/sdk/io/kafka/KafkaIO.java | 156 +++-
.../KafkaIOReadImplementationCompatibility.java | 6 +
.../beam/sdk/io/kafka/ReadFromKafkaDoFn.java | 114 +--
.../io/messaging-expansion-service/build.gradle | 3 +
.../io/snowflake/crosslanguage/package-info.java | 4 -
.../beam/sdk/io/solace/broker/BrokerResponse.java | 13 +-
.../sdk/io/solace/broker/BrokerResponseTest.java | 66 ++
sdks/java/io/sparkreceiver/3/build.gradle | 6 +
.../apache/beam/sdk/testing/kafka/LocalKafka.java | 2 +-
sdks/python/apache_beam/coders/row_coder_test.py | 54 ++
sdks/python/apache_beam/dataframe/io.py | 2 +-
sdks/python/apache_beam/dataframe/io_test.py | 15 +-
.../anomaly_detection_pipeline/setup.py | 2 +-
.../online_clustering/clustering_pipeline/setup.py | 2 +-
...ytorch_image_classification_with_side_inputs.py | 2 +-
.../transforms/elementwise/enrichment_test.py | 2 +-
.../apache_beam/io/external/xlang_jmsio_it_test.py | 383 ++++++++++
.../io/external/xlang_mqttio_it_test.py | 66 +-
sdks/python/apache_beam/io/filesystemio.py | 5 +-
sdks/python/apache_beam/io/gcp/__init__.py | 20 -
.../apache_beam/io/gcp/bigquery_change_history.py | 10 +-
.../apache_beam/io/gcp/bigquery_file_loads.py | 83 +-
.../apache_beam/io/gcp/bigquery_file_loads_test.py | 151 ++++
sdks/python/apache_beam/io/gcp/bigquery_test.py | 4 +-
.../apache_beam/io/gcp/bigquery_write_it_test.py | 71 +-
.../apache_beam/io/gcp/bigtableio_it_test.py | 10 +-
.../apache_beam/io/gcp/gcsfilesystem_test.py | 4 +-
sdks/python/apache_beam/io/watch.py | 737 ++++++++++++++++++
sdks/python/apache_beam/io/watch_test.py | 475 ++++++++++++
sdks/python/apache_beam/ml/anomaly/specifiable.py | 2 +-
.../apache_beam/ml/anomaly/univariate/mean_test.py | 8 +-
.../apache_beam/ml/anomaly/univariate/perf_test.py | 51 +-
.../ml/anomaly/univariate/quantile_test.py | 8 +-
.../ml/anomaly/univariate/stdev_test.py | 8 +-
sdks/python/apache_beam/ml/inference/base.py | 5 +-
.../python/apache_beam/options/pipeline_options.py | 5 +
.../apache_beam/options/pipeline_options_test.py | 25 +
sdks/python/apache_beam/portability/common_urns.py | 1 +
sdks/python/apache_beam/runners/common.pxd | 1 +
sdks/python/apache_beam/runners/common.py | 35 +
sdks/python/apache_beam/runners/common_test.py | 58 ++
.../runners/dataflow/internal/clients/README.txt | 11 -
.../runners/dataflow/internal/clients/__init__.py | 16 -
.../apache_beam/runners/dataflow/internal/names.py | 2 +-
.../dataproc/dataproc_cluster_manager.py | 3 +-
.../runners/interactive/interactive_beam_test.py | 8 +-
.../runners/interactive/recording_manager.py | 197 +++--
.../runners/interactive/recording_manager_test.py | 191 ++++-
.../runners/portability/beam_plugins_it_test.py | 70 ++
.../runners/portability/prism_runner.py | 5 +-
.../apache_beam/runners/worker/sdk_worker_main.py | 20 +-
.../runners/worker/sdk_worker_main_test.py | 49 ++
sdks/python/apache_beam/testing/util.py | 43 ++
sdks/python/apache_beam/testing/util_test.py | 44 ++
.../apache_beam/transforms/async_dofn_test.py | 15 +-
.../transforms/managed_iceberg_it_test.py | 4 +-
sdks/python/apache_beam/typehints/schemas.py | 94 ++-
sdks/python/apache_beam/typehints/schemas_test.py | 94 +++
sdks/python/apache_beam/typehints/typehints.py | 13 +-
.../python/apache_beam/typehints/typehints_test.py | 18 +
sdks/python/apache_beam/utils/timestamp.py | 327 ++++++--
sdks/python/apache_beam/utils/timestamp_test.py | 187 ++++-
.../databases/debezium.yaml} | 41 +-
sdks/python/apache_beam/yaml/integration_tests.py | 73 +-
sdks/python/apache_beam/yaml/standard_io.yaml | 23 +
sdks/python/build.gradle | 4 +-
.../container/base_image_requirements_manual.txt | 1 +
sdks/python/container/boot.go | 74 +-
.../container/ml/py310/base_image_requirements.txt | 1 +
.../container/ml/py310/gpu_image_requirements.txt | 1 +
.../container/ml/py311/base_image_requirements.txt | 1 +
.../container/ml/py311/gpu_image_requirements.txt | 1 +
.../container/ml/py312/base_image_requirements.txt | 1 +
.../container/ml/py312/gpu_image_requirements.txt | 1 +
.../container/ml/py313/base_image_requirements.txt | 1 +
sdks/python/container/piputil.go | 39 +-
sdks/python/container/profiler.go | 303 +++++++-
sdks/python/container/profiler_test.go | 125 +++
.../container/py310/base_image_requirements.txt | 1 +
.../container/py311/base_image_requirements.txt | 1 +
.../container/py312/base_image_requirements.txt | 1 +
.../container/py313/base_image_requirements.txt | 1 +
.../container/py314/base_image_requirements.txt | 1 +
sdks/python/pyproject.toml | 111 ++-
sdks/python/scripts/run_integration_test.sh | 11 +-
sdks/python/scripts/run_lint.sh | 37 +-
sdks/python/setup.py | 32 +-
sdks/python/test-suites/dataflow/common.gradle | 12 +-
sdks/python/test-suites/direct/common.gradle | 21 +-
sdks/python/tox.ini | 11 +-
sdks/standard_expansion_services.yaml | 4 +
sdks/standard_external_transforms.yaml | 103 ++-
settings.gradle.kts | 1 +
.../site/content/en/documentation/io/connectors.md | 34 +-
.../site/content/en/documentation/io/managed-io.md | 68 ++
.../site/content/en/documentation/runners/flink.md | 23 +-
.../sdks/python-pipeline-dependencies.md | 41 +-
374 files changed, 14764 insertions(+), 2543 deletions(-)
copy .github/trigger_files/{beam_PostCommit_Java_ValidatesRunner_Spark.json =>
beam_PostCommit_Python_Arm.json} (100%)
create mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/MultiKeyBundleOptions.java
copy
sdks/java/core/src/main/java/org/apache/beam/sdk/values/OpenTelemetryContextPropagator.java
=>
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillOpenTelemetryContextPropagator.java
(76%)
rename
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/{KeyTokenInvalidException.java
=> WorkCancellingException.java} (50%)
create mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounterTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/KeyTokenInvalidExceptionTest.java
create mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateTest.java
create mode 100644
sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/BeamSqlEnvRegisterOperatorTest.java
copy sdks/java/io/{synthetic => arrow-flight}/build.gradle (65%)
create mode 100644
sdks/java/io/arrow-flight/src/main/java/org/apache/beam/sdk/io/arrowflight/ArrowFlightIO.java
copy
sdks/java/{extensions/sql/udf/src/main/java/org/apache/beam/sdk/extensions/sql/udf
=>
io/arrow-flight/src/main/java/org/apache/beam/sdk/io/arrowflight}/package-info.java
(69%)
create mode 100644
sdks/java/io/arrow-flight/src/test/java/org/apache/beam/sdk/io/arrowflight/ArrowFlightIOTest.java
create mode 100644
sdks/java/io/debezium/src/test/java/org/apache/beam/io/debezium/DebeziumReadSchemaTransformProviderTest.java
create mode 100644
sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/CreateCDCReadTasksDoFn.java
create mode 100644
sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaCDCReadTask.java
create mode 100644
sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaCDCSourceDoFn.java
copy
sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/{AutoScaler.java =>
BeamGenericJmsConnectionFactory.java} (52%)
create mode 100644
sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/ConnectionConfiguration.java
copy
sdks/java/io/{mqtt/src/main/java/org/apache/beam/sdk/io/mqtt/MqttReadSchemaTransformProvider.java
=>
jms/src/main/java/org/apache/beam/sdk/io/jms/JmsReadSchemaTransformProvider.java}
(50%)
create mode 100644
sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/JmsWriteSchemaTransformProvider.java
create mode 100644
sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/ConnectionConfigurationTest.java
create mode 100644
sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/JmsSchemaTransformProviderTest.java
create mode 100644
sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/broker/BrokerResponseTest.java
create mode 100644 sdks/python/apache_beam/io/external/xlang_jmsio_it_test.py
create mode 100644 sdks/python/apache_beam/io/watch.py
create mode 100644 sdks/python/apache_beam/io/watch_test.py
delete mode 100644
sdks/python/apache_beam/runners/dataflow/internal/clients/README.txt
delete mode 100644
sdks/python/apache_beam/runners/dataflow/internal/clients/__init__.py
create mode 100644
sdks/python/apache_beam/runners/portability/beam_plugins_it_test.py
copy sdks/python/apache_beam/yaml/{tests/map.yaml =>
extended_tests/databases/debezium.yaml} (55%)
create mode 100644 sdks/python/container/profiler_test.go