This is an automated email from the ASF dual-hosted git repository. github-bot pushed a change to tag nightly-master in repository https://gitbox.apache.org/repos/asf/beam.git.
*** WARNING: tag nightly-master was modified! *** from 3615801 (commit) to 5fb31eb (commit) from 3615801 [BEAM-12334] Re-use java 11 flag in build.gradle (#14892) add edad245 [BEAM-12380] Add KafkaIO Transforms and Kafka Taxi example. add 6e0ce6d [BEAM-12380] Go Kafka Fixup, documenting Taxi example, and (de)serializer TODO. add 0908b24 Merge pull request #14996: [BEAM-12380] Add KafkaIO Transforms and Kafka Taxi example. add 3bca2e3 [BEAM-12454] Add original filesToStage to final pipeline options add acc0f76 Merge pull request #14946 from kw2542/BEAM-12454 add a1950ef Bump version of Cloud Datastore Python package (#15017) add c5a6988 [BEAM-9487] Disable GBK safety checks by default add f2bbae8 Merge pull request #15003 from zhoufek/aut add defbc1b [BEAM-12297] Add methods to PubsubIO for reading DynamicMessage add be7717f Merge pull request #14971: [BEAM-12297] Add methods to PubsubIO for reading DynamicMessage add 6a1c932 Enable Portable Job Submission for Runner V2 for large graphs. add 2088b66 Merge pull request #14724: Enable Portable Job Submission for Runner V2 for large graphs. add d63d2ee [BEAM-11833] Minor code clarity improvements. add 9f34aab [BEAM-11833] Minor code clarity improvements. add 5fb31eb [BEAM-8137] Refactor External Worker Service to be part of sdks-java-harness (#14988) No new revisions were added by this update. Summary of changes: CHANGES.md | 3 + .../beam/runners/flink/FlinkJobServerDriver.java | 2 +- .../worker/BeamFnMapTaskExecutorFactory.java | 2 +- .../worker/DataflowMapTaskExecutorFactory.java | 2 +- .../dataflow/worker/DataflowRunnerHarness.java | 4 +- .../worker/IntrinsicMapTaskExecutorFactory.java | 2 +- .../dataflow/worker/SdkHarnessRegistries.java | 2 +- .../dataflow/worker/SdkHarnessRegistry.java | 2 +- .../dataflow/worker/fn/BeamFnControlService.java | 2 +- .../worker/fn/data/BeamFnDataGrpcService.java | 2 +- .../worker/fn/logging/BeamFnLoggingService.java | 2 +- .../worker/fn/BeamFnControlServiceTest.java | 4 +- .../worker/fn/data/BeamFnDataGrpcServiceTest.java | 2 +- .../fn/logging/BeamFnLoggingServiceTest.java | 2 +- .../artifact/ArtifactRetrievalService.java | 2 +- .../artifact/ArtifactStagingService.java | 2 +- .../control/DefaultJobBundleFactory.java | 6 +- .../control/FnApiControlClientPoolService.java | 4 +- .../SingleEnvironmentInstanceJobBundleFactory.java | 2 +- .../runners/fnexecution/data/GrpcDataService.java | 2 +- .../environment/DockerEnvironmentFactory.java | 4 +- .../environment/EmbeddedEnvironmentFactory.java | 6 +- .../environment/EnvironmentFactory.java | 4 +- .../environment/ExternalEnvironmentFactory.java | 4 +- .../environment/ProcessEnvironmentFactory.java | 2 +- .../StaticRemoteEnvironmentFactory.java | 2 +- .../fnexecution/logging/GrpcLoggingService.java | 2 +- .../provisioning/StaticGrpcProvisionService.java | 4 +- .../fnexecution/state/GrpcStateService.java | 2 +- .../status/BeamWorkerStatusGrpcService.java | 4 +- .../runners/fnexecution/EmbeddedSdkHarness.java | 3 + .../GrpcContextHeaderAccessorProviderTest.java | 2 + .../runners/fnexecution/ServerFactoryTest.java | 1 + .../artifact/ArtifactRetrievalServiceTest.java | 2 +- .../control/DefaultJobBundleFactoryTest.java | 4 +- .../control/FnApiControlClientPoolServiceTest.java | 6 +- .../fnexecution/control/RemoteExecutionTest.java | 6 +- ...gleEnvironmentInstanceJobBundleFactoryTest.java | 4 +- .../fnexecution/data/GrpcDataServiceTest.java | 4 +- .../environment/DockerEnvironmentFactoryTest.java | 2 +- .../environment/ProcessEnvironmentFactoryTest.java | 2 +- .../logging/GrpcLoggingServiceTest.java | 4 +- .../StaticGrpcProvisionServiceTest.java | 6 +- .../status/BeamWorkerStatusGrpcServiceTest.java | 6 +- .../runners/jobsubmission/InMemoryJobService.java | 4 +- .../runners/jobsubmission/JobServerDriver.java | 4 +- .../jobsubmission/InMemoryJobServiceTest.java | 2 +- runners/portability/java/build.gradle | 1 - .../beam/runners/portability/PortableRunner.java | 5 +- .../beam/runners/samza/SamzaJobServerDriver.java | 2 +- .../beam/runners/spark/SparkJobServerDriver.java | 2 +- sdks/go/examples/kafka/taxi.go | 163 ++++++++++++ sdks/go/pkg/beam/core/graph/edge.go | 7 +- sdks/go/pkg/beam/core/runtime/graphx/xlang.go | 114 ++++---- sdks/go/pkg/beam/io/pubsubio/pubsubio.go | 2 +- sdks/go/pkg/beam/io/xlang/kafkaio/kafka.go | 288 +++++++++++++++++++++ .../src/main/java/org/apache/beam/sdk/io/Read.java | 15 +- .../org/apache/beam/sdk/fn/server}/FnService.java | 2 +- .../server}/GrpcContextHeaderAccessorProvider.java | 2 +- .../apache/beam/sdk/fn/server}/GrpcFnServer.java | 2 +- .../apache/beam/sdk/fn/server}/HeaderAccessor.java | 2 +- .../sdk/fn/server}/InProcessServerFactory.java | 2 +- .../apache/beam/sdk/fn/server}/ServerFactory.java | 2 +- .../apache/beam/sdk/fn/server}/package-info.java | 4 +- .../beam/fn/harness}/ExternalWorkerService.java | 9 +- .../fn/harness}/ExternalWorkerServiceTest.java | 2 +- sdks/java/io/expansion-service/build.gradle | 6 + sdks/java/io/google-cloud-platform/build.gradle | 1 + .../apache/beam/sdk/io/gcp/pubsub/PubsubIO.java | 53 ++++ .../beam/sdk/io/gcp/pubsub/PubsubIOTest.java | 97 +++++++ .../python/apache_beam/options/pipeline_options.py | 3 +- .../runners/dataflow/internal/apiclient.py | 4 + sdks/python/apache_beam/transforms/core.py | 39 ++- .../apache_beam/transforms/ptransform_test.py | 8 +- sdks/python/apache_beam/transforms/trigger.py | 6 +- sdks/python/container/base_image_requirements.txt | 2 +- 76 files changed, 814 insertions(+), 178 deletions(-) create mode 100644 sdks/go/examples/kafka/taxi.go create mode 100644 sdks/go/pkg/beam/io/xlang/kafkaio/kafka.go rename {runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution => sdks/java/fn-execution/src/main/java/org/apache/beam/sdk/fn/server}/FnService.java (97%) rename {runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution => sdks/java/fn-execution/src/main/java/org/apache/beam/sdk/fn/server}/GrpcContextHeaderAccessorProvider.java (98%) rename {runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution => sdks/java/fn-execution/src/main/java/org/apache/beam/sdk/fn/server}/GrpcFnServer.java (99%) rename {runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution => sdks/java/fn-execution/src/main/java/org/apache/beam/sdk/fn/server}/HeaderAccessor.java (95%) rename {runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution => sdks/java/fn-execution/src/main/java/org/apache/beam/sdk/fn/server}/InProcessServerFactory.java (98%) rename {runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution => sdks/java/fn-execution/src/main/java/org/apache/beam/sdk/fn/server}/ServerFactory.java (99%) copy sdks/java/{testing/tpcds/src/main/java/org/apache/beam/sdk/tpcds => fn-execution/src/main/java/org/apache/beam/sdk/fn/server}/package-info.java (92%) rename {runners/portability/java/src/main/java/org/apache/beam/runners/portability => sdks/java/harness/src/main/java/org/apache/beam/fn/harness}/ExternalWorkerService.java (95%) rename {runners/portability/java/src/test/java/org/apache/beam/runners/portability => sdks/java/harness/src/test/java/org/apache/beam/fn/harness}/ExternalWorkerServiceTest.java (98%)