See
<https://builds.apache.org/job/beam_PostCommit_Go_VR_Flink/2664/display/redirect?page=changes>
Changes:
[git] Remove optionality and add sensible defaults to PubsubIO builders.
------------------------------------------
[...truncated 814.06 KB...]
[CHAIN MapPartition (MapPartition at [3]{github.com, Flatten}) -> FlatMap
(FlatMap at ExtractOutput[0]) (1/1)] INFO
org.apache.flink.runtime.taskmanager.Task - Freeing task resources for CHAIN
MapPartition (MapPartition at [3]{github.com, Flatten}) -> FlatMap (FlatMap at
ExtractOutput[0]) (1/1) (36f7c649c5a62a9592dcdd24ab42d5fd).
[CHAIN MapPartition (MapPartition at [3]{github.com, Flatten}) -> FlatMap
(FlatMap at ExtractOutput[0]) (1/1)] INFO
org.apache.flink.runtime.taskmanager.Task - Ensuring all FileSystem streams are
closed for task CHAIN MapPartition (MapPartition at [3]{github.com, Flatten})
-> FlatMap (FlatMap at ExtractOutput[0]) (1/1)
(36f7c649c5a62a9592dcdd24ab42d5fd) [FINISHED]
[flink-akka.actor.default-dispatcher-5] INFO
org.apache.flink.runtime.taskexecutor.TaskExecutor - Un-registering task and
sending final execution state FINISHED to JobManager for task CHAIN
MapPartition (MapPartition at [3]{github.com, Flatten}) -> FlatMap (FlatMap at
ExtractOutput[0]) 36f7c649c5a62a9592dcdd24ab42d5fd.
[flink-akka.actor.default-dispatcher-4] INFO
org.apache.flink.runtime.executiongraph.ExecutionGraph - GroupReduce
(GroupReduce at CoGBK) (1/1) (83aa8a231c1bc8068cea60ab0862de34) switched from
CREATED to SCHEDULED.
[flink-akka.actor.default-dispatcher-4] INFO
org.apache.flink.runtime.executiongraph.ExecutionGraph - GroupReduce
(GroupReduce at CoGBK) (1/1) (83aa8a231c1bc8068cea60ab0862de34) switched from
SCHEDULED to DEPLOYING.
[flink-akka.actor.default-dispatcher-4] INFO
org.apache.flink.runtime.executiongraph.ExecutionGraph - Deploying GroupReduce
(GroupReduce at CoGBK) (1/1) (attempt #0) to
b9aa5d5d-4e54-4525-871f-56e4a4703821 @ localhost (dataPort=-1)
[flink-akka.actor.default-dispatcher-5] INFO
org.apache.flink.runtime.taskexecutor.TaskExecutor - Received task GroupReduce
(GroupReduce at CoGBK) (1/1).
[flink-akka.actor.default-dispatcher-4] INFO
org.apache.flink.runtime.executiongraph.ExecutionGraph - CHAIN MapPartition
(MapPartition at [3]{github.com, Flatten}) -> FlatMap (FlatMap at
ExtractOutput[0]) (1/1) (36f7c649c5a62a9592dcdd24ab42d5fd) switched from
RUNNING to FINISHED.
[GroupReduce (GroupReduce at CoGBK) (1/1)] INFO
org.apache.flink.runtime.taskmanager.Task - GroupReduce (GroupReduce at CoGBK)
(1/1) (83aa8a231c1bc8068cea60ab0862de34) switched from CREATED to DEPLOYING.
[GroupReduce (GroupReduce at CoGBK) (1/1)] INFO
org.apache.flink.runtime.taskmanager.Task - Creating FileSystem stream leak
safety net for task GroupReduce (GroupReduce at CoGBK) (1/1)
(83aa8a231c1bc8068cea60ab0862de34) [DEPLOYING]
[GroupReduce (GroupReduce at CoGBK) (1/1)] INFO
org.apache.flink.runtime.taskmanager.Task - Loading JAR files for task
GroupReduce (GroupReduce at CoGBK) (1/1) (83aa8a231c1bc8068cea60ab0862de34)
[DEPLOYING].
[GroupReduce (GroupReduce at CoGBK) (1/1)] INFO
org.apache.flink.runtime.taskmanager.Task - Registering task at network:
GroupReduce (GroupReduce at CoGBK) (1/1) (83aa8a231c1bc8068cea60ab0862de34)
[DEPLOYING].
[GroupReduce (GroupReduce at CoGBK) (1/1)] INFO
org.apache.flink.runtime.taskmanager.Task - GroupReduce (GroupReduce at CoGBK)
(1/1) (83aa8a231c1bc8068cea60ab0862de34) switched from DEPLOYING to RUNNING.
[flink-akka.actor.default-dispatcher-6] INFO
org.apache.flink.runtime.executiongraph.ExecutionGraph - GroupReduce
(GroupReduce at CoGBK) (1/1) (83aa8a231c1bc8068cea60ab0862de34) switched from
DEPLOYING to RUNNING.
[CHAIN Filter (UnionFixFilter) -> Map (Key Extractor) -> GroupCombine
(GroupCombine at GroupCombine: CoGBK) -> Map (Key Extractor) (1/1)] INFO
org.apache.flink.runtime.taskmanager.Task - CHAIN Filter (UnionFixFilter) ->
Map (Key Extractor) -> GroupCombine (GroupCombine at GroupCombine: CoGBK) ->
Map (Key Extractor) (1/1) (f7561365ae76fa3149e9ceeca0180e86) switched from
RUNNING to FINISHED.
[CHAIN Filter (UnionFixFilter) -> Map (Key Extractor) -> GroupCombine
(GroupCombine at GroupCombine: CoGBK) -> Map (Key Extractor) (1/1)] INFO
org.apache.flink.runtime.taskmanager.Task - Freeing task resources for CHAIN
Filter (UnionFixFilter) -> Map (Key Extractor) -> GroupCombine (GroupCombine at
GroupCombine: CoGBK) -> Map (Key Extractor) (1/1)
(f7561365ae76fa3149e9ceeca0180e86).
[CHAIN Filter (UnionFixFilter) -> Map (Key Extractor) -> GroupCombine
(GroupCombine at GroupCombine: CoGBK) -> Map (Key Extractor) (1/1)] INFO
org.apache.flink.runtime.taskmanager.Task - Ensuring all FileSystem streams are
closed for task CHAIN Filter (UnionFixFilter) -> Map (Key Extractor) ->
GroupCombine (GroupCombine at GroupCombine: CoGBK) -> Map (Key Extractor) (1/1)
(f7561365ae76fa3149e9ceeca0180e86) [FINISHED]
[flink-akka.actor.default-dispatcher-4] INFO
org.apache.flink.runtime.taskexecutor.TaskExecutor - Un-registering task and
sending final execution state FINISHED to JobManager for task CHAIN Filter
(UnionFixFilter) -> Map (Key Extractor) -> GroupCombine (GroupCombine at
GroupCombine: CoGBK) -> Map (Key Extractor) f7561365ae76fa3149e9ceeca0180e86.
[flink-akka.actor.default-dispatcher-4] INFO
org.apache.flink.runtime.executiongraph.ExecutionGraph - CHAIN Filter
(UnionFixFilter) -> Map (Key Extractor) -> GroupCombine (GroupCombine at
GroupCombine: CoGBK) -> Map (Key Extractor) (1/1)
(f7561365ae76fa3149e9ceeca0180e86) switched from RUNNING to FINISHED.
[flink-akka.actor.default-dispatcher-4] INFO
org.apache.flink.runtime.executiongraph.ExecutionGraph - MapPartition
(MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
(1/1) (0161884c0972fb5c938f9e09e09fbcc2) switched from CREATED to SCHEDULED.
[flink-akka.actor.default-dispatcher-4] INFO
org.apache.flink.runtime.executiongraph.ExecutionGraph - MapPartition
(MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
(1/1) (0161884c0972fb5c938f9e09e09fbcc2) switched from SCHEDULED to DEPLOYING.
[flink-akka.actor.default-dispatcher-4] INFO
org.apache.flink.runtime.executiongraph.ExecutionGraph - Deploying MapPartition
(MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
(1/1) (attempt #0) to b9aa5d5d-4e54-4525-871f-56e4a4703821 @ localhost
(dataPort=-1)
[flink-akka.actor.default-dispatcher-3] INFO
org.apache.flink.runtime.taskexecutor.TaskExecutor - Received task MapPartition
(MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
(1/1).
[MapPartition (MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
(1/1)] INFO org.apache.flink.runtime.taskmanager.Task - MapPartition
(MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
(1/1) (0161884c0972fb5c938f9e09e09fbcc2) switched from CREATED to DEPLOYING.
[MapPartition (MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
(1/1)] INFO org.apache.flink.runtime.taskmanager.Task - Creating FileSystem
stream leak safety net for task MapPartition (MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
(1/1) (0161884c0972fb5c938f9e09e09fbcc2) [DEPLOYING]
[MapPartition (MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
(1/1)] INFO org.apache.flink.runtime.taskmanager.Task - Loading JAR files for
task MapPartition (MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
(1/1) (0161884c0972fb5c938f9e09e09fbcc2) [DEPLOYING].
[MapPartition (MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
(1/1)] INFO org.apache.flink.runtime.taskmanager.Task - Registering task at
network: MapPartition (MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
(1/1) (0161884c0972fb5c938f9e09e09fbcc2) [DEPLOYING].
[MapPartition (MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
(1/1)] INFO org.apache.flink.runtime.taskmanager.Task - MapPartition
(MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
(1/1) (0161884c0972fb5c938f9e09e09fbcc2) switched from DEPLOYING to RUNNING.
[flink-akka.actor.default-dispatcher-3] INFO
org.apache.flink.runtime.executiongraph.ExecutionGraph - MapPartition
(MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
(1/1) (0161884c0972fb5c938f9e09e09fbcc2) switched from DEPLOYING to RUNNING.
[MapPartition (MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
(1/1)] WARN org.apache.flink.metrics.MetricGroup - The operator name
MapPartition (MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
exceeded the 80 characters length limit and was truncated.
[GroupReduce (GroupReduce at CoGBK) (1/1)] INFO
org.apache.flink.runtime.taskmanager.Task - GroupReduce (GroupReduce at CoGBK)
(1/1) (83aa8a231c1bc8068cea60ab0862de34) switched from RUNNING to FINISHED.
[GroupReduce (GroupReduce at CoGBK) (1/1)] INFO
org.apache.flink.runtime.taskmanager.Task - Freeing task resources for
GroupReduce (GroupReduce at CoGBK) (1/1) (83aa8a231c1bc8068cea60ab0862de34).
[GroupReduce (GroupReduce at CoGBK) (1/1)] INFO
org.apache.flink.runtime.taskmanager.Task - Ensuring all FileSystem streams are
closed for task GroupReduce (GroupReduce at CoGBK) (1/1)
(83aa8a231c1bc8068cea60ab0862de34) [FINISHED]
[flink-akka.actor.default-dispatcher-7] INFO
org.apache.flink.runtime.taskexecutor.TaskExecutor - Un-registering task and
sending final execution state FINISHED to JobManager for task GroupReduce
(GroupReduce at CoGBK) 83aa8a231c1bc8068cea60ab0862de34.
[flink-akka.actor.default-dispatcher-7] INFO
org.apache.flink.runtime.executiongraph.ExecutionGraph - GroupReduce
(GroupReduce at CoGBK) (1/1) (83aa8a231c1bc8068cea60ab0862de34) switched from
RUNNING to FINISHED.
[grpc-default-executor-5] INFO
org.apache.beam.runners.fnexecution.artifact.AbstractArtifactRetrievalService -
GetManifest for
/tmp/beam-artifact-staging/go-job-7-1583441936923625831_e051ac32-68e7-4506-af8a-98fe91df044d/MANIFEST
[grpc-default-executor-5] INFO
org.apache.beam.runners.fnexecution.artifact.AbstractArtifactRetrievalService -
GetManifest for
/tmp/beam-artifact-staging/go-job-7-1583441936923625831_e051ac32-68e7-4506-af8a-98fe91df044d/MANIFEST
-> 1 artifacts
[grpc-default-executor-4] INFO
org.apache.beam.runners.fnexecution.control.FnApiControlClientPoolService -
Beam Fn Control client connected with id 18-1
[grpc-default-executor-5] INFO
org.apache.beam.runners.fnexecution.logging.GrpcLoggingService - Beam Fn
Logging client connected.
[grpc-default-executor-5] INFO
<https://builds.apache.org/job/beam_PostCommit_Go_VR_Flink/ws/src/sdks/go/test/.gogradle/project_gopath/src/github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/core/runtime/harness/harness.go>:331
- Connecting via grpc @ localhost:44271 ...
[grpc-default-executor-5] INFO
<https://builds.apache.org/job/beam_PostCommit_Go_VR_Flink/ws/src/sdks/go/test/.gogradle/project_gopath/src/github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/core/runtime/harness/harness.go>:331
- Connecting via grpc @ localhost:43971 ...
[grpc-default-executor-4] INFO
<https://builds.apache.org/job/beam_PostCommit_Go_VR_Flink/ws/src/sdks/go/test/.gogradle/project_gopath/src/github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/core/runtime/harness/harness.go>:331
- Connecting via grpc @ localhost:34013 ...
[grpc-default-executor-4] INFO
org.apache.beam.runners.fnexecution.data.GrpcDataService - Beam Fn Data client
connected.
[grpc-default-executor-4] INFO
<https://builds.apache.org/job/beam_PostCommit_Go_VR_Flink/ws/src/sdks/go/test/.gogradle/project_gopath/src/github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/core/runtime/exec/datasource.go>:246
- DataSource: 1 elements in 3145522 ns
[MapPartition (MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
(1/1)] INFO
org.apache.beam.runners.fnexecution.control.DefaultJobBundleFactory - Closing
environment urn: "beam:env:docker:v1"
payload: "\nAus.gcr.io/apache-beam-testing/jenkins/beam_go_sdk:20200305-205339"
capabilities: "beam:protocol:progress_reporting:v0"
capabilities: "beam:protocol:multi_core_bundle_processing:v1"
capabilities: "beam:coder:bytes:v1"
capabilities: "beam:coder:bool:v1"
capabilities: "beam:coder:varint:v1"
capabilities: "beam:coder:double:v1"
capabilities: "beam:coder:length_prefix:v1"
capabilities: "beam:coder:kv:v1"
capabilities: "beam:coder:iterable:v1"
capabilities: "beam:coder:state_backed_iterable:v1"
capabilities: "beam:coder:windowed_value:v1"
capabilities: "beam:coder:global_window:v1"
capabilities: "beam:coder:interval_window:v1"
[MapPartition (MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
(1/1)] INFO org.apache.beam.runners.fnexecution.logging.GrpcLoggingService - 1
Beam Fn Logging clients still connected during shutdown.
[MapPartition (MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
(1/1)] WARN org.apache.beam.sdk.fn.data.BeamFnDataGrpcMultiplexer - Hanged up
for unknown endpoint.
[MapPartition (MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
(1/1)] INFO
org.apache.beam.runners.fnexecution.environment.DockerContainerEnvironment -
Closing Docker container
1d2575746880f623c6b07dc21421d47f761550a03216d68d8be6dee62df82ca5. Logs:
2020/03/05 20:59:05 Provision info:
pipeline_options:<fields:<key:"beam:option:app_name:v1"
value:<string_value:"go-job-7-1583441936923625831" > >
fields:<key:"beam:option:experiments:v1"
value:<list_value:<values:<string_value:"beam_fn_api" > > > >
fields:<key:"beam:option:flink_master:v1" value:<string_value:"[local]" > >
fields:<key:"beam:option:go_options:v1"
value:<struct_value:<fields:<key:"options"
value:<struct_value:<fields:<key:"hooks" value:<string_value:"{}" > > > > > > >
> fields:<key:"beam:option:job_name:v1"
value:<string_value:"go0job0701583441936923625831-jenkins-0305205857-9f66694f"
> > fields:<key:"beam:option:options_id:v1" value:<number_value:7 > >
fields:<key:"beam:option:output_executable_path:v1"
value:<null_value:NULL_VALUE > > fields:<key:"beam:option:runner:v1"
value:<null_value:NULL_VALUE > > >
retrieval_token:"/tmp/beam-artifact-staging/go-job-7-1583441936923625831_e051ac32-68e7-4506-af8a-98fe91df044d/MANIFEST"
logging_endpoint:<url:"localhost:43971" >
artifact_endpoint:<url:"localhost:46799" >
control_endpoint:<url:"localhost:44271" >
2020/03/05 20:59:05 Initializing Go harness: /opt/apache/beam/boot --id=18-1
--provision_endpoint=localhost:41723
Worker exited successfully!
Failed to send message: rpc error: code = Unavailable desc = transport is
closing
severity:WARN timestamp:<seconds:1583441946 nanos:176234127 > message:"forcing
DataChannel[localhost:34013] reconnection on port {localhost:34013} due to rpc
error: code = Canceled desc = Multiplexer hanging up" instruction_id:"2"
log_location:"<https://builds.apache.org/job/beam_PostCommit_Go_VR_Flink/ws/src/sdks/go/test/.gogradle/project_gopath/src/github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/core/runtime/harness/datamgr.go>:119"
Remote logging failed: rpc error: code = Unavailable desc = transport is
closing. Retrying in 5 sec ...
[MapPartition (MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
(1/1)] WARN
org.apache.beam.runners.fnexecution.control.DefaultJobBundleFactory - Error
cleaning up servers urn: "beam:env:docker:v1"
payload: "\nAus.gcr.io/apache-beam-testing/jenkins/beam_go_sdk:20200305-205339"
capabilities: "beam:protocol:progress_reporting:v0"
capabilities: "beam:protocol:multi_core_bundle_processing:v1"
capabilities: "beam:coder:bytes:v1"
capabilities: "beam:coder:bool:v1"
capabilities: "beam:coder:varint:v1"
capabilities: "beam:coder:double:v1"
capabilities: "beam:coder:length_prefix:v1"
capabilities: "beam:coder:kv:v1"
capabilities: "beam:coder:iterable:v1"
capabilities: "beam:coder:state_backed_iterable:v1"
capabilities: "beam:coder:windowed_value:v1"
capabilities: "beam:coder:global_window:v1"
capabilities: "beam:coder:interval_window:v1"
java.io.IOException: Received exit code 1 for command 'docker rm
1d2575746880f623c6b07dc21421d47f761550a03216d68d8be6dee62df82ca5'. stderr:
Error: No such container:
1d2575746880f623c6b07dc21421d47f761550a03216d68d8be6dee62df82ca5
at
org.apache.beam.runners.fnexecution.environment.DockerCommand.runShortCommand(DockerCommand.java:234)
at
org.apache.beam.runners.fnexecution.environment.DockerCommand.runShortCommand(DockerCommand.java:168)
at
org.apache.beam.runners.fnexecution.environment.DockerCommand.removeContainer(DockerCommand.java:163)
at
org.apache.beam.runners.fnexecution.environment.DockerContainerEnvironment.close(DockerContainerEnvironment.java:95)
at
org.apache.beam.runners.fnexecution.control.DefaultJobBundleFactory$WrappedSdkHarnessClient.$closeResource(DefaultJobBundleFactory.java:479)
at
org.apache.beam.runners.fnexecution.control.DefaultJobBundleFactory$WrappedSdkHarnessClient.close(DefaultJobBundleFactory.java:479)
at
org.apache.beam.runners.fnexecution.control.DefaultJobBundleFactory$WrappedSdkHarnessClient.unref(DefaultJobBundleFactory.java:494)
at
org.apache.beam.runners.fnexecution.control.DefaultJobBundleFactory$WrappedSdkHarnessClient.access$1600(DefaultJobBundleFactory.java:432)
at
org.apache.beam.runners.fnexecution.control.DefaultJobBundleFactory.lambda$createEnvironmentCaches$3(DefaultJobBundleFactory.java:169)
at
org.apache.beam.vendor.guava.v26_0_jre.com.google.common.cache.LocalCache.processPendingNotifications(LocalCache.java:1809)
at
org.apache.beam.vendor.guava.v26_0_jre.com.google.common.cache.LocalCache$Segment.runUnlockedCleanup(LocalCache.java:3462)
at
org.apache.beam.vendor.guava.v26_0_jre.com.google.common.cache.LocalCache$Segment.postWriteCleanup(LocalCache.java:3438)
at
org.apache.beam.vendor.guava.v26_0_jre.com.google.common.cache.LocalCache$Segment.clear(LocalCache.java:3215)
at
org.apache.beam.vendor.guava.v26_0_jre.com.google.common.cache.LocalCache.clear(LocalCache.java:4270)
at
org.apache.beam.vendor.guava.v26_0_jre.com.google.common.cache.LocalCache$LocalManualCache.invalidateAll(LocalCache.java:4909)
at
org.apache.beam.runners.fnexecution.control.DefaultJobBundleFactory.close(DefaultJobBundleFactory.java:259)
at
org.apache.beam.runners.fnexecution.control.DefaultExecutableStageContext.close(DefaultExecutableStageContext.java:43)
at
org.apache.beam.runners.fnexecution.control.ReferenceCountingExecutableStageContextFactory$WrappedContext.closeActual(ReferenceCountingExecutableStageContextFactory.java:208)
at
org.apache.beam.runners.fnexecution.control.ReferenceCountingExecutableStageContextFactory$WrappedContext.access$200(ReferenceCountingExecutableStageContextFactory.java:184)
at
org.apache.beam.runners.fnexecution.control.ReferenceCountingExecutableStageContextFactory.release(ReferenceCountingExecutableStageContextFactory.java:173)
at
org.apache.beam.runners.fnexecution.control.ReferenceCountingExecutableStageContextFactory.scheduleRelease(ReferenceCountingExecutableStageContextFactory.java:132)
at
org.apache.beam.runners.fnexecution.control.ReferenceCountingExecutableStageContextFactory.access$300(ReferenceCountingExecutableStageContextFactory.java:44)
at
org.apache.beam.runners.fnexecution.control.ReferenceCountingExecutableStageContextFactory$WrappedContext.close(ReferenceCountingExecutableStageContextFactory.java:204)
at
org.apache.beam.runners.flink.translation.functions.FlinkExecutableStageFunction.$closeResource(FlinkExecutableStageFunction.java:204)
at
org.apache.beam.runners.flink.translation.functions.FlinkExecutableStageFunction.close(FlinkExecutableStageFunction.java:291)
at
org.apache.flink.api.common.functions.util.FunctionUtils.closeFunction(FunctionUtils.java:43)
at org.apache.flink.runtime.operators.BatchTask.run(BatchTask.java:508)
at
org.apache.flink.runtime.operators.BatchTask.invoke(BatchTask.java:369)
at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:705)
at org.apache.flink.runtime.taskmanager.Task.run(Task.java:530)
at java.lang.Thread.run(Thread.java:748)
[MapPartition (MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
(1/1)] INFO org.apache.flink.runtime.taskmanager.Task - MapPartition
(MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
(1/1) (0161884c0972fb5c938f9e09e09fbcc2) switched from RUNNING to FINISHED.
[MapPartition (MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
(1/1)] INFO org.apache.flink.runtime.taskmanager.Task - Freeing task resources
for MapPartition (MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
(1/1) (0161884c0972fb5c938f9e09e09fbcc2).
[flink-akka.actor.default-dispatcher-7] INFO
org.apache.flink.runtime.executiongraph.ExecutionGraph - DataSink
(DiscardingOutput) (1/1) (1366ca2e1b9948fecdf935d5e66c278a) switched from
CREATED to SCHEDULED.
[MapPartition (MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
(1/1)] INFO org.apache.flink.runtime.taskmanager.Task - Ensuring all
FileSystem streams are closed for task MapPartition (MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
(1/1) (0161884c0972fb5c938f9e09e09fbcc2) [FINISHED]
[flink-akka.actor.default-dispatcher-7] INFO
org.apache.flink.runtime.executiongraph.ExecutionGraph - DataSink
(DiscardingOutput) (1/1) (1366ca2e1b9948fecdf935d5e66c278a) switched from
SCHEDULED to DEPLOYING.
[flink-akka.actor.default-dispatcher-3] INFO
org.apache.flink.runtime.taskexecutor.TaskExecutor - Un-registering task and
sending final execution state FINISHED to JobManager for task MapPartition
(MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
0161884c0972fb5c938f9e09e09fbcc2.
[flink-akka.actor.default-dispatcher-7] INFO
org.apache.flink.runtime.executiongraph.ExecutionGraph - Deploying DataSink
(DiscardingOutput) (1/1) (attempt #0) to b9aa5d5d-4e54-4525-871f-56e4a4703821 @
localhost (dataPort=-1)
[flink-akka.actor.default-dispatcher-7] INFO
org.apache.flink.runtime.executiongraph.ExecutionGraph - MapPartition
(MapPartition at
[1]github.com/apache/beam/sdks/go/test/vendor/github.com/apache/beam/sdks/go/pkg/beam/testing/passert.sumFn)
(1/1) (0161884c0972fb5c938f9e09e09fbcc2) switched from RUNNING to FINISHED.
[flink-akka.actor.default-dispatcher-3] INFO
org.apache.flink.runtime.taskexecutor.TaskExecutor - Received task DataSink
(DiscardingOutput) (1/1).
[DataSink (DiscardingOutput) (1/1)] INFO
org.apache.flink.runtime.taskmanager.Task - DataSink (DiscardingOutput) (1/1)
(1366ca2e1b9948fecdf935d5e66c278a) switched from CREATED to DEPLOYING.
[DataSink (DiscardingOutput) (1/1)] INFO
org.apache.flink.runtime.taskmanager.Task - Creating FileSystem stream leak
safety net for task DataSink (DiscardingOutput) (1/1)
(1366ca2e1b9948fecdf935d5e66c278a) [DEPLOYING]
[DataSink (DiscardingOutput) (1/1)] INFO
org.apache.flink.runtime.taskmanager.Task - Loading JAR files for task DataSink
(DiscardingOutput) (1/1) (1366ca2e1b9948fecdf935d5e66c278a) [DEPLOYING].
[DataSink (DiscardingOutput) (1/1)] INFO
org.apache.flink.runtime.taskmanager.Task - Registering task at network:
DataSink (DiscardingOutput) (1/1) (1366ca2e1b9948fecdf935d5e66c278a)
[DEPLOYING].
[DataSink (DiscardingOutput) (1/1)] INFO
org.apache.flink.runtime.taskmanager.Task - DataSink (DiscardingOutput) (1/1)
(1366ca2e1b9948fecdf935d5e66c278a) switched from DEPLOYING to RUNNING.
[flink-akka.actor.default-dispatcher-3] INFO
org.apache.flink.runtime.executiongraph.ExecutionGraph - DataSink
(DiscardingOutput) (1/1) (1366ca2e1b9948fecdf935d5e66c278a) switched from
DEPLOYING to RUNNING.
[DataSink (DiscardingOutput) (1/1)] INFO
org.apache.flink.runtime.taskmanager.Task - DataSink (DiscardingOutput) (1/1)
(1366ca2e1b9948fecdf935d5e66c278a) switched from RUNNING to FINISHED.
[DataSink (DiscardingOutput) (1/1)] INFO
org.apache.flink.runtime.taskmanager.Task - Freeing task resources for DataSink
(DiscardingOutput) (1/1) (1366ca2e1b9948fecdf935d5e66c278a).
[DataSink (DiscardingOutput) (1/1)] INFO
org.apache.flink.runtime.taskmanager.Task - Ensuring all FileSystem streams are
closed for task DataSink (DiscardingOutput) (1/1)
(1366ca2e1b9948fecdf935d5e66c278a) [FINISHED]
[flink-akka.actor.default-dispatcher-3] INFO
org.apache.flink.runtime.taskexecutor.TaskExecutor - Un-registering task and
sending final execution state FINISHED to JobManager for task DataSink
(DiscardingOutput) 1366ca2e1b9948fecdf935d5e66c278a.
[flink-akka.actor.default-dispatcher-3] INFO
org.apache.flink.runtime.executiongraph.ExecutionGraph - DataSink
(DiscardingOutput) (1/1) (1366ca2e1b9948fecdf935d5e66c278a) switched from
RUNNING to FINISHED.
[flink-akka.actor.default-dispatcher-3] INFO
org.apache.flink.runtime.executiongraph.ExecutionGraph - Job
go0job0701583441936923625831-jenkins-0305205857-9f66694f
(2df28e4811fc5b6bd35afe128e309587) switched from state RUNNING to FINISHED.
[flink-akka.actor.default-dispatcher-3] INFO
org.apache.flink.runtime.dispatcher.StandaloneDispatcher - Job
2df28e4811fc5b6bd35afe128e309587 reached globally terminal state FINISHED.
[flink-akka.actor.default-dispatcher-8] INFO
org.apache.flink.runtime.jobmaster.JobMaster - Stopping the JobMaster for job
go0job0701583441936923625831-jenkins-0305205857-9f66694f(2df28e4811fc5b6bd35afe128e309587).
[flink-akka.actor.default-dispatcher-8] INFO
org.apache.flink.runtime.jobmaster.slotpool.SlotPoolImpl - Suspending SlotPool.
[flink-akka.actor.default-dispatcher-8] INFO
org.apache.flink.runtime.jobmaster.JobMaster - Close ResourceManager connection
24781fd6767db8fb3616359737ee34c9: JobManager is shutting down..
[flink-akka.actor.default-dispatcher-8] INFO
org.apache.flink.runtime.jobmaster.slotpool.SlotPoolImpl - Stopping SlotPool.
[flink-akka.actor.default-dispatcher-5] INFO
org.apache.flink.runtime.resourcemanager.StandaloneResourceManager - Disconnect
job manager 900076656bdb38212b95807f7d904a36@akka://flink/user/jobmanager_13
for job 2df28e4811fc5b6bd35afe128e309587 from the resource manager.
[flink-akka.actor.default-dispatcher-3] INFO
org.apache.flink.runtime.taskexecutor.slot.TaskSlotTable - Free slot
TaskSlot(index:0, state:ACTIVE, resource profile:
ResourceProfile{cpuCores=1.7976931348623157E308, heapMemoryInMB=2147483647,
directMemoryInMB=2147483647, nativeMemoryInMB=2147483647,
networkMemoryInMB=2147483647, managedMemoryInMB=16273}, allocationId:
1923f019370096ef8f4d9df5e6b55760, jobId: 2df28e4811fc5b6bd35afe128e309587).
[flink-akka.actor.default-dispatcher-3] INFO
org.apache.flink.runtime.taskexecutor.JobLeaderService - Remove job
2df28e4811fc5b6bd35afe128e309587 from job leader monitoring.
[flink-runner-job-invoker] INFO
org.apache.flink.runtime.minicluster.MiniCluster - Shutting down Flink Mini
Cluster
[flink-akka.actor.default-dispatcher-3] INFO
org.apache.flink.runtime.taskexecutor.TaskExecutor - Close JobManager
connection for job 2df28e4811fc5b6bd35afe128e309587.
[flink-runner-job-invoker] INFO
org.apache.flink.runtime.dispatcher.DispatcherRestEndpoint - Shutting down rest
endpoint.
[flink-akka.actor.default-dispatcher-3] INFO
org.apache.flink.runtime.taskexecutor.TaskExecutor - Close JobManager
connection for job 2df28e4811fc5b6bd35afe128e309587.
[flink-akka.actor.default-dispatcher-3] INFO
org.apache.flink.runtime.taskexecutor.JobLeaderService - Cannot reconnect to
job 2df28e4811fc5b6bd35afe128e309587 because it is not registered.
[flink-akka.actor.default-dispatcher-3] INFO
org.apache.flink.runtime.taskexecutor.TaskExecutor - Stopping TaskExecutor
akka://flink/user/taskmanager_12.
[flink-akka.actor.default-dispatcher-3] INFO
org.apache.flink.runtime.taskexecutor.TaskExecutor - Close ResourceManager
connection 24781fd6767db8fb3616359737ee34c9.
[flink-akka.actor.default-dispatcher-3] INFO
org.apache.flink.runtime.taskexecutor.JobLeaderService - Stop job leader
service.
[flink-akka.actor.default-dispatcher-7] INFO
org.apache.flink.runtime.resourcemanager.StandaloneResourceManager - Closing
TaskExecutor connection b9aa5d5d-4e54-4525-871f-56e4a4703821 because: The
TaskExecutor is shutting down.
[flink-akka.actor.default-dispatcher-3] INFO
org.apache.flink.runtime.state.TaskExecutorLocalStateStoresManager - Shutting
down TaskExecutorLocalStateStoresManager.
[flink-akka.actor.default-dispatcher-3] INFO
org.apache.flink.runtime.io.disk.FileChannelManagerImpl - FileChannelManager
removed spill file directory /tmp/flink-io-413e218d-57ed-4699-b763-67758762c159
[flink-akka.actor.default-dispatcher-3] INFO
org.apache.flink.runtime.io.network.NettyShuffleEnvironment - Shutting down the
network environment and its components.
[flink-akka.actor.default-dispatcher-3] INFO
org.apache.flink.runtime.io.disk.FileChannelManagerImpl - FileChannelManager
removed spill file directory
/tmp/flink-netty-shuffle-8b1e20cf-74bc-44de-98d5-2c86a7c84c32
[flink-akka.actor.default-dispatcher-3] INFO
org.apache.flink.runtime.taskexecutor.KvStateService - Shutting down the
kvState service and its components.
[flink-akka.actor.default-dispatcher-3] INFO
org.apache.flink.runtime.taskexecutor.JobLeaderService - Stop job leader
service.
[flink-akka.actor.default-dispatcher-3] INFO
org.apache.flink.runtime.filecache.FileCache - removed file cache directory
/tmp/flink-dist-cache-ffa94fd8-82f1-4fd3-b842-b2b454d46138
[flink-akka.actor.default-dispatcher-3] INFO
org.apache.flink.runtime.taskexecutor.TaskExecutor - Stopped TaskExecutor
akka://flink/user/taskmanager_12.
[ForkJoinPool.commonPool-worker-4] INFO
org.apache.flink.runtime.dispatcher.DispatcherRestEndpoint - Removing cache
directory /tmp/flink-web-ui
[ForkJoinPool.commonPool-worker-4] INFO
org.apache.flink.runtime.dispatcher.DispatcherRestEndpoint - Shut down complete.
[flink-akka.actor.default-dispatcher-3] INFO
org.apache.flink.runtime.resourcemanager.StandaloneResourceManager - Shut down
cluster because application is in CANCELED, diagnostics
DispatcherResourceManagerComponent has been closed..
[flink-akka.actor.default-dispatcher-7] INFO
org.apache.flink.runtime.dispatcher.StandaloneDispatcher - Stopping dispatcher
akka://flink/user/dispatcher.
[flink-akka.actor.default-dispatcher-8] INFO
org.apache.flink.runtime.resourcemanager.slotmanager.SlotManagerImpl - Closing
the SlotManager.
[flink-akka.actor.default-dispatcher-7] INFO
org.apache.flink.runtime.dispatcher.StandaloneDispatcher - Stopping all
currently running jobs of dispatcher akka://flink/user/dispatcher.
[flink-akka.actor.default-dispatcher-8] INFO
org.apache.flink.runtime.resourcemanager.slotmanager.SlotManagerImpl -
Suspending the SlotManager.
[flink-akka.actor.default-dispatcher-7] INFO
org.apache.flink.runtime.rest.handler.legacy.backpressure.StackTraceSampleCoordinator
- Shutting down stack trace sample coordinator.
[flink-akka.actor.default-dispatcher-7] INFO
org.apache.flink.runtime.dispatcher.StandaloneDispatcher - Stopped dispatcher
akka://flink/user/dispatcher.
[flink-akka.actor.default-dispatcher-7] INFO
org.apache.flink.runtime.rpc.akka.AkkaRpcService - Stopping Akka RPC service.
[flink-metrics-2] INFO akka.remote.RemoteActorRefProvider$RemotingTerminator -
Shutting down remote daemon.
[flink-metrics-2] INFO akka.remote.RemoteActorRefProvider$RemotingTerminator -
Remote daemon shut down; proceeding with flushing remote transports.
[flink-metrics-2] INFO akka.remote.RemoteActorRefProvider$RemotingTerminator -
Remoting shut down.
[flink-metrics-2] INFO org.apache.flink.runtime.rpc.akka.AkkaRpcService -
Stopping Akka RPC service.
[flink-metrics-2] INFO org.apache.flink.runtime.rpc.akka.AkkaRpcService -
Stopped Akka RPC service.
[flink-akka.actor.default-dispatcher-8] INFO
org.apache.flink.runtime.blob.PermanentBlobCache - Shutting down BLOB cache
[flink-akka.actor.default-dispatcher-8] INFO
org.apache.flink.runtime.blob.TransientBlobCache - Shutting down BLOB cache
[flink-akka.actor.default-dispatcher-8] INFO
org.apache.flink.runtime.blob.BlobServer - Stopped BLOB server at 0.0.0.0:40517
[flink-akka.actor.default-dispatcher-8] INFO
org.apache.flink.runtime.rpc.akka.AkkaRpcService - Stopped Akka RPC service.
[flink-runner-job-invoker] INFO
org.apache.beam.runners.flink.FlinkPipelineRunner - Execution finished in 8199
msecs
[flink-runner-job-invoker] INFO
org.apache.beam.runners.flink.FlinkPipelineRunner - Final accumulator values:
[flink-runner-job-invoker] INFO
org.apache.beam.runners.flink.FlinkPipelineRunner - __metricscontainers :
MetricQueryResults(Counters(n1/beam:env:docker:v1:0:beam:metric:element_count:v1
{PCOLLECTION=n8:0}: 3, n5/beam:env:docker:v1:0:beam:metric:element_count:v1
{PCOLLECTION=n5}: 1, n3/beam:env:docker:v1:0:beam:metric:element_count:v1
{PCOLLECTION=n3}: 1, n5/beam:env:docker:v1:0:beam:metric:element_count:v1
{PCOLLECTION=n7}: 3, n3/beam:env:docker:v1:0:beam:metric:element_count:v1
{PCOLLECTION=n4}: 3, n3/beam:env:docker:v1:0:beam:metric:element_count:v1
{PCOLLECTION=n7}: 3, n1/beam:env:docker:v1:0:beam:metric:element_count:v1
{PCOLLECTION=n2}: 3, n1/beam:env:docker:v1:0:beam:metric:element_count:v1
{PCOLLECTION=n7}: 3, n5/beam:env:docker:v1:0:beam:metric:element_count:v1
{PCOLLECTION=n6}: 3, n1/beam:env:docker:v1:0:beam:metric:element_count:v1
{PCOLLECTION=n1}: 1, n3/beam:env:docker:v1:0:beam:metric:element_count:v1
{PCOLLECTION=n8:1}: 3, n9/beam:env:docker:v1:0:beam:metric:element_count:v1
{PCOLLECTION=n9}: 1, n5/beam:env:docker:v1:0:beam:metric:element_count:v1
{PCOLLECTION=n8:2}:
3)Distributions(n1/beam:env:docker:v1:0:beam:metric:sampled_byte_size:v1
{PCOLLECTION=n2}: DistributionResult{sum=4, count=2, min=2, max=2},
n3/beam:env:docker:v1:0:beam:metric:sampled_byte_size:v1 {PCOLLECTION=n7}:
DistributionResult{sum=4, count=2, min=2, max=2},
n5/beam:env:docker:v1:0:beam:metric:sampled_byte_size:v1 {PCOLLECTION=n7}:
DistributionResult{sum=2, count=1, min=2, max=2},
n1/beam:env:docker:v1:0:beam:metric:sampled_byte_size:v1 {PCOLLECTION=n1}:
DistributionResult{sum=1, count=1, min=1, max=1},
n3/beam:env:docker:v1:0:beam:metric:sampled_byte_size:v1 {PCOLLECTION=n8:1}:
DistributionResult{sum=8, count=2, min=4, max=4},
n5/beam:env:docker:v1:0:beam:metric:sampled_byte_size:v1 {PCOLLECTION=n5}:
DistributionResult{sum=1, count=1, min=1, max=1},
n9/beam:env:docker:v1:0:beam:metric:sampled_byte_size:v1 {PCOLLECTION=n9}:
DistributionResult{sum=20, count=1, min=20, max=20},
n5/beam:env:docker:v1:0:beam:metric:sampled_byte_size:v1 {PCOLLECTION=n8:2}:
DistributionResult{sum=12, count=3, min=4, max=4},
n3/beam:env:docker:v1:0:beam:metric:sampled_byte_size:v1 {PCOLLECTION=n4}:
DistributionResult{sum=4, count=2, min=2, max=2},
n3/beam:env:docker:v1:0:beam:metric:sampled_byte_size:v1 {PCOLLECTION=n3}:
DistributionResult{sum=1, count=1, min=1, max=1}))
[flink-runner-job-invoker] INFO
org.apache.beam.runners.fnexecution.artifact.AbstractArtifactRetrievalService -
Manifest at
/tmp/beam-artifact-staging/go-job-7-1583441936923625831_e051ac32-68e7-4506-af8a-98fe91df044d/MANIFEST
has 1 artifact locations
[flink-runner-job-invoker] INFO
org.apache.beam.runners.fnexecution.artifact.BeamFileSystemArtifactStagingService
- Removed dir
/tmp/beam-artifact-staging/go-job-7-1583441936923625831_e051ac32-68e7-4506-af8a-98fe91df044d/
2020/03/05 20:59:07 Job state: DONE
2020/03/05 20:59:07 Test flatten:flatten completed
2020/03/05 20:59:07 Result: 1 tests failed
if [[ ! -z "$JOB_PORT" ]]; then
# Shut down the job server
kill %1 || echo "Failed to shut down job server"
fi
# Delete the container locally and remotely
docker rmi $CONTAINER:$TAG || echo "Failed to remove container"
Error response from daemon: conflict: unable to remove repository reference
"us.gcr.io/apache-beam-testing/jenkins/beam_go_sdk:20200305-205339" (must
force) - container 8c5404b263c8 is using its referenced image 66502089c971
Failed to remove container
gcloud --quiet container images delete $CONTAINER:$TAG || echo "Failed to
delete container"
Digests:
-
us.gcr.io/apache-beam-testing/jenkins/beam_go_sdk@sha256:4e6ca8a9bfbd35548fa8833b23a09fe6f493e3756ae48f8720b66f6140bc4e1f
Associated tags:
- 20200305-205339
Tags:
- us.gcr.io/apache-beam-testing/jenkins/beam_go_sdk:20200305-205339
Deleted [us.gcr.io/apache-beam-testing/jenkins/beam_go_sdk:20200305-205339].
Deleted
[us.gcr.io/apache-beam-testing/jenkins/beam_go_sdk@sha256:4e6ca8a9bfbd35548fa8833b23a09fe6f493e3756ae48f8720b66f6140bc4e1f].
# Clean up tempdir
rm -rf $TMPDIR
if [[ "$TEST_EXIT_CODE" -eq 0 ]]; then
echo ">>> SUCCESS"
else
echo ">>> FAILURE"
fi
exit $TEST_EXIT_CODE
>>> FAILURE
> Task :sdks:go:test:flinkValidatesRunner FAILED
FAILURE: Build failed with an exception.
* Where:
Build file
'<https://builds.apache.org/job/beam_PostCommit_Go_VR_Flink/ws/src/sdks/go/test/build.gradle'>
line: 59
* What went wrong:
Execution failed for task ':sdks:go:test:flinkValidatesRunner'.
> Process 'command 'sh'' finished with non-zero exit value 1
* Try:
Run with --stacktrace option to get the stack trace. Run with --info or --debug
option to get more log output. Run with --scan to get full insights.
* Get more help at https://help.gradle.org
Deprecated Gradle features were used in this build, making it incompatible with
Gradle 6.0.
Use '--warning-mode all' to show the individual deprecation warnings.
See
https://docs.gradle.org/5.2.1/userguide/command_line_interface.html#sec:command_line_warnings
BUILD FAILED in 7m 39s
67 actionable tasks: 50 executed, 17 from cache
Publishing build scan...
https://gradle.com/s/3bteh6j7a7ybm
Build step 'Invoke Gradle script' changed build result to FAILURE
Build step 'Invoke Gradle script' marked build as failure
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]