HJK nomad
05/16/2023, 6:08 AMenv = StreamExecutionEnvironment.get_execution_environment()
env.set_runtime_mode(RuntimeExecutionMode.BATCH)
env.set_parallelism(1)
my_array = [1, 2, 3, 4, 5]
stream = env.from_collection(collection=my_array)
stream.print()
env.execute()
but job finished, and tm log error.
2023-05-16 054832,411 INFO org.apache.beam.runners.fnexecution.logging.GrpcLoggingService [] - 1 Beam Fn Logging clients still connected during shutdown.
2023-05-16 054832,421 WARN org.apache.beam.sdk.fn.data.BeamFnDataGrpcMultiplexer [] - Hanged up for unknown endpoint.
2023-05-16 054832,426 WARN org.apache.beam.sdk.fn.data.BeamFnDataGrpcMultiplexer [] - Hanged up for unknown endpoint.
2023-05-16 054832,487 WARN org.apache.beam.runners.fnexecution.control.DefaultJobBundleFactory [] - Error cleaning up servers urn: "beamenvprocess:v1"
payload: "\032H/usr/local/lib/python3.8/dist-packages/pyflink/bin/pyflink-udf-runner.sh\"\225\002\n\004PATH\022\214\002/root/miniconda3/condabin:/usr/liDian Fu
05/16/2023, 7:48 AMHJK nomad
05/17/2023, 1:10 AMHJK nomad
05/17/2023, 1:21 AMHJK nomad
05/18/2023, 1:42 AM