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%)