Using Apache Beam & Flink as a Runner, Executing Simple Word Count Program. If Input File is Greater than 1GB, Observing below Exception..
Tried With Smaller Size Files Around ~900MB File and WordCount Program Executed Successfully
root@ravi:~# python3 wordcount_with_metrics.py --input beam/bigfile.txt --out bigfile-output.txt
INFO:apache_beam.runners.worker.worker_pool_main:Listening for workers at localhost:38423
WARNING:root:Make sure that locally built Python SDK docker image has Python 3.8 interpreter.
INFO:root:Default Python SDK image for environment is apache/beam_python3.8_sdk:2.29.0
INFO:apache_beam.runners.portability.fn_api_runner.translations:==================== <function lift_combiners at 0x7fdec4d209d0> ====================
INFO:apache_beam.runners.portability.fn_api_runner.translations:==================== <function sort_stages at 0x7fdec4d21160> ====================
INFO:apache_beam.runners.portability.flink_runner:Adding HTTP protocol scheme to flink_master parameter: http://localhost:8081
INFO:apache_beam.utils.subprocess_server:Using cached job server jar from https://repo.maven.apache.org/maven2/org/apache/beam/beam-runners-flink-1.12-job-server/2.29.0/beam-runners-flink-1.12-job-server-2.29.0.jar
INFO:apache_beam.utils.subprocess_server:Starting service with ['java' '-jar' '/home/root/.apache_beam/cache/jars/beam-runners-flink-1.12-job-server-2.29.0.jar' '--flink-master' 'http://localhost:8081' '--artifacts-dir' '/tmp/beam-temp5isb3pa3/artifactsy6mkh_yt' '--job-port' '34737' '--artifact-port' '0' '--expansion-port' '0']
INFO:apache_beam.utils.subprocess_server:b'Jun 03, 2021 5:33:23 AM org.apache.beam.runners.jobsubmission.JobServerDriver createArtifactStagingService'
INFO:apache_beam.utils.subprocess_server:b'INFO: ArtifactStagingService started on localhost:33829'
INFO:apache_beam.utils.subprocess_server:b'Jun 03, 2021 5:33:24 AM org.apache.beam.runners.jobsubmission.JobServerDriver createExpansionService'
INFO:apache_beam.utils.subprocess_server:b'INFO: Java ExpansionService started on localhost:35453'
INFO:apache_beam.utils.subprocess_server:b'Jun 03, 2021 5:33:24 AM org.apache.beam.runners.jobsubmission.JobServerDriver createJobServer'
INFO:apache_beam.utils.subprocess_server:b'INFO: JobService started on localhost:34737'
INFO:apache_beam.utils.subprocess_server:b'Jun 03, 2021 5:33:24 AM org.apache.beam.runners.jobsubmission.JobServerDriver run'
INFO:apache_beam.utils.subprocess_server:b'INFO: Job server now running, terminate with Ctrl+C'
WARNING:apache_beam.options.pipeline_options:Discarding unparseable args: [[]]
INFO:apache_beam.utils.subprocess_server:b'Jun 03, 2021 5:33:25 AM org.apache.beam.runners.fnexecution.artifact.ArtifactStagingService$2 onNext'
INFO:apache_beam.utils.subprocess_server:b'INFO: Staging artifacts for job_9d7e5a97-a367-4892-a685-02c1e0cf3d5a.'
INFO:apache_beam.utils.subprocess_server:b'Jun 03, 2021 5:33:25 AM org.apache.beam.runners.fnexecution.artifact.ArtifactStagingService$2 resolveNextEnvironment'
INFO:apache_beam.utils.subprocess_server:b'INFO: Resolving artifacts for job_9d7e5a97-a367-4892-a685-02c1e0cf3d5a.ref_Environment_default_environment_1.'
INFO:apache_beam.utils.subprocess_server:b'Jun 03, 2021 5:33:25 AM org.apache.beam.runners.fnexecution.artifact.ArtifactStagingService$2 onNext'
INFO:apache_beam.utils.subprocess_server:b'INFO: Getting 1 artifacts for job_9d7e5a97-a367-4892-a685-02c1e0cf3d5a.null.'
INFO:apache_beam.utils.subprocess_server:b'Jun 03, 2021 5:33:25 AM org.apache.beam.runners.fnexecution.artifact.ArtifactStagingService$2 finishStaging'
INFO:apache_beam.utils.subprocess_server:b'INFO: Artifacts fully staged for job_9d7e5a97-a367-4892-a685-02c1e0cf3d5a.'
INFO:apache_beam.utils.subprocess_server:b'Jun 03, 2021 5:33:26 AM org.apache.beam.runners.flink.FlinkJobInvoker invokeWithExecutor'
INFO:apache_beam.utils.subprocess_server:b'INFO: Invoking job BeamApp-root-0603053325-adad87a2_db245f01-11b8-4ac2-b74a-efd73a9f9b71 with pipeline runner org.apache.beam.runners.flink.FlinkPipelineRunner@473c0e7b'
INFO:apache_beam.utils.subprocess_server:b'Jun 03, 2021 5:33:26 AM org.apache.beam.runners.jobsubmission.JobInvocation start'
INFO:apache_beam.utils.subprocess_server:b'INFO: Starting job invocation BeamApp-root-0603053325-adad87a2_db245f01-11b8-4ac2-b74a-efd73a9f9b71'
INFO:apache_beam.runners.portability.portable_runner:Environment "LOOPBACK" has started a component necessary for the execution. Be sure to run the pipeline using
with Pipeline() as p:
p.apply(..)
This ensures that the pipeline finishes before this program exits.
INFO:apache_beam.runners.portability.portable_runner:Job state changed to STOPPED
INFO:apache_beam.runners.portability.portable_runner:Job state changed to STARTING
INFO:apache_beam.runners.portability.portable_runner:Job state changed to RUNNING
INFO:apache_beam.utils.subprocess_server:b'Jun 03, 2021 5:33:26 AM org.apache.beam.runners.flink.FlinkPipelineRunner runPipelineWithTranslator'
INFO:apache_beam.utils.subprocess_server:b'INFO: Translating pipeline to Flink program.'
INFO:apache_beam.utils.subprocess_server:b'Jun 03, 2021 5:33:26 AM org.apache.beam.runners.flink.FlinkExecutionEnvironments createBatchExecutionEnvironment'
INFO:apache_beam.utils.subprocess_server:b'INFO: Creating a Batch Execution Environment.'
INFO:apache_beam.utils.subprocess_server:b'Jun 03, 2021 5:33:26 AM org.apache.beam.runners.flink.FlinkExecutionEnvironments createBatchExecutionEnvironment'
INFO:apache_beam.utils.subprocess_server:b'INFO: Using Flink Master URL localhost:8081.'
INFO:apache_beam.utils.subprocess_server:b'Jun 03, 2021 5:33:26 AM org.apache.flink.api.java.utils.PlanGenerator logTypeRegistrationDetails'
INFO:apache_beam.utils.subprocess_server:b'INFO: The job has 0 registered types and 0 default Kryo serializers'
INFO:apache_beam.runners.worker.statecache:Creating state cache with size 0
INFO:apache_beam.runners.worker.sdk_worker:Creating insecure control channel for localhost:41679.
INFO:apache_beam.runners.worker.sdk_worker:Control channel established.
INFO:apache_beam.runners.worker.sdk_worker:Initializing SDKHarness with unbounded number of workers.
INFO:apache_beam.runners.worker.sdk_worker:Creating insecure state channel for localhost:34167.
INFO:apache_beam.runners.worker.sdk_worker:State channel established.
INFO:apache_beam.runners.worker.data_plane:Creating client data channel for localhost:35057
INFO:apache_beam.runners.worker.sdk_worker:No more requests from control plane
INFO:apache_beam.runners.worker.sdk_worker:SDK Harness waiting for in-flight requests to complete
INFO:apache_beam.runners.worker.data_plane:Closing all cached grpc data channels.
INFO:apache_beam.runners.worker.sdk_worker:Closing all cached gRPC state handlers.
INFO:apache_beam.runners.worker.sdk_worker:Done consuming work.
INFO:apache_beam.runners.worker.statecache:Creating state cache with size 0
INFO:apache_beam.runners.worker.sdk_worker:Creating insecure control channel for localhost:33243.
INFO:apache_beam.runners.worker.sdk_worker:Control channel established.
INFO:apache_beam.runners.worker.sdk_worker:Initializing SDKHarness with unbounded number of workers.
INFO:apache_beam.runners.worker.sdk_worker:Creating insecure state channel for localhost:46869.
INFO:apache_beam.runners.worker.sdk_worker:State channel established.
INFO:apache_beam.runners.worker.data_plane:Creating client data channel for localhost:39177
INFO:apache_beam.runners.worker.sdk_worker:No more requests from control plane
INFO:apache_beam.runners.worker.sdk_worker:SDK Harness waiting for in-flight requests to complete
INFO:apache_beam.runners.worker.data_plane:Closing all cached grpc data channels.
INFO:apache_beam.runners.worker.sdk_worker:Closing all cached gRPC state handlers.
INFO:apache_beam.runners.worker.sdk_worker:Done consuming work.
INFO:apache_beam.runners.worker.statecache:Creating state cache with size 0
INFO:apache_beam.runners.worker.sdk_worker:Creating insecure control channel for localhost:40859.
INFO:apache_beam.runners.worker.sdk_worker:Control channel established.
INFO:apache_beam.runners.worker.sdk_worker:Initializing SDKHarness with unbounded number of workers.
INFO:apache_beam.runners.worker.sdk_worker:Creating insecure state channel for localhost:37567.
INFO:apache_beam.runners.worker.sdk_worker:State channel established.
INFO:apache_beam.runners.worker.data_plane:Creating client data channel for localhost:36685
E0603 05:45:57.654085397 262 chttp2_transport.cc:1117] Received a GOAWAY with error code ENHANCE_YOUR_CALM and debug data equal to "too_many_pings"
ERROR:apache_beam.runners.worker.data_plane:Failed to read inputs in the data plane.
Traceback (most recent call last):
File "/usr/local/lib/python3.8/dist-packages/apache_beam/runners/worker/data_plane.py", line 581, in _read_inputs
for elements in elements_iterator:
File "/usr/local/lib/python3.8/dist-packages/grpc/_channel.py", line 426, in __next__
return self._next()
File "/usr/local/lib/python3.8/dist-packages/grpc/_channel.py", line 826, in _next
raise self
grpc._channel._MultiThreadedRendezvous: <_MultiThreadedRendezvous of RPC that terminated with:
status = StatusCode.UNAVAILABLE
details = "Socket closed"
debug_error_string = "{"created":"@1622699157.655022787","description":"Error received from peer ipv6:[::1]:36685","file":"src/core/lib/surface/call.cc","file_line":1066,"grpc_message":"Socket closed","grpc_status":14}"
>
Exception in thread read_grpc_client_inputs:
Traceback (most recent call last):
File "/usr/lib/python3.8/threading.py", line 932, in _bootstrap_inner
self.run()
File "/usr/lib/python3.8/threading.py", line 870, in run
self._target(*self._args, **self._kwargs)
File "/usr/local/lib/python3.8/dist-packages/apache_beam/runners/worker/data_plane.py", line 598, in <lambda>
target=lambda: self._read_inputs(elements_iterator),
File "/usr/local/lib/python3.8/dist-packages/apache_beam/runners/worker/data_plane.py", line 581, in _read_inputs
for elements in elements_iterator:
File "/usr/local/lib/python3.8/dist-packages/grpc/_channel.py", line 426, in __next__
return self._next()
File "/usr/local/lib/python3.8/dist-packages/grpc/_channel.py", line 826, in _next
raise self
grpc._channel._MultiThreadedRendezvous: <_MultiThreadedRendezvous of RPC that terminated with:
status = StatusCode.UNAVAILABLE
details = "Socket closed"
debug_error_string = "{"created":"@1622699157.655022787","description":"Error received from peer ipv6:[::1]:36685","file":"src/core/lib/surface/call.cc","file_line":1066,"grpc_message":"Socket closed","grpc_status":14}"
>
Traceback (most recent call last):
File "/usr/local/lib/python3.8/dist-packages/apache_beam/runners/worker/data_plane.py", line 470, in input_elements
element = received.get(timeout=1)
File "/usr/lib/python3.8/queue.py", line 178, in get
raise Empty
_queue.Empty
During handling of the above exception, another exception occurred:
Traceback (most recent call last):
File "/usr/local/lib/python3.8/dist-packages/apache_beam/runners/worker/sdk_worker.py", line 289, in _execute
response = task()
File "/usr/local/lib/python3.8/dist-packages/apache_beam/runners/worker/sdk_worker.py", line 362, in <lambda>
lambda: self.create_worker().do_instruction(request), request)
File "/usr/local/lib/python3.8/dist-packages/apache_beam/runners/worker/sdk_worker.py", line 606, in do_instruction
return getattr(self, request_type)(
File "/usr/local/lib/python3.8/dist-packages/apache_beam/runners/worker/sdk_worker.py", line 644, in process_bundle
bundle_processor.process_bundle(instruction_id))
File "/usr/local/lib/python3.8/dist-packages/apache_beam/runners/worker/bundle_processor.py", line 989, in process_bundle
for element in data_channel.input_elements(instruction_id,
File "/usr/local/lib/python3.8/dist-packages/apache_beam/runners/worker/data_plane.py", line 473, in input_elements
raise RuntimeError('Channel closed prematurely.')
RuntimeError: Channel closed prematurely.
ERROR:apache_beam.runners.worker.sdk_worker:Error processing instruction 5. Original traceback is
Traceback (most recent call last):
File "/usr/local/lib/python3.8/dist-packages/apache_beam/runners/worker/data_plane.py", line 470, in input_elements
element = received.get(timeout=1)
File "/usr/lib/python3.8/queue.py", line 178, in get
raise Empty
_queue.Empty
During handling of the above exception, another exception occurred:
Traceback (most recent call last):
File "/usr/local/lib/python3.8/dist-packages/apache_beam/runners/worker/sdk_worker.py", line 289, in _execute
response = task()
File "/usr/local/lib/python3.8/dist-packages/apache_beam/runners/worker/sdk_worker.py", line 362, in <lambda>
lambda: self.create_worker().do_instruction(request), request)
File "/usr/local/lib/python3.8/dist-packages/apache_beam/runners/worker/sdk_worker.py", line 606, in do_instruction
return getattr(self, request_type)(
File "/usr/local/lib/python3.8/dist-packages/apache_beam/runners/worker/sdk_worker.py", line 644, in process_bundle
bundle_processor.process_bundle(instruction_id))
File "/usr/local/lib/python3.8/dist-packages/apache_beam/runners/worker/bundle_processor.py", line 989, in process_bundle
for element in data_channel.input_elements(instruction_id,
File "/usr/local/lib/python3.8/dist-packages/apache_beam/runners/worker/data_plane.py", line 473, in input_elements
raise RuntimeError('Channel closed prematurely.')
RuntimeError: Channel closed prematurely.
- Beam: 2.31.0
- Flink 1.12.0
- Python 3.8.5
- Java 1.8.0_282
To Generate Big File
- tr -dc "A-Za-z 0-9" < /dev/urandom | fold -w100|head -n 15000000 > bigfile.txt