Hi,
I'm trying to read parquet file with Flink 1.12.0 Scala API and save it as another parquet file. Now it's working correctly with ParquetRowInputFormat: val inputPath: String = ... val messageType: MessageType = ... val parquetInputFormat = new ParquetRowInputFormat(new Path(inputPath), messageType) parquetInputFormat.setNestedFileEnumeration(true) env.readFile(parquetInputFormat, inputPath) .map(row => {....//mapping row to MyPOJO}) .sinkTo(FileSink.forBulkFormat...) But when I replace the inputFormat: val pojoTypeInfo = Types.POJO(classOf[MyPOJO]).asInstanceOf[PojoTypeInfo[MyPOJO]] val parquetInputFormat = new ParquetPojoInputFormat(new Path(inputPath), messageType, pojoTypeInfo) parquetInputFormat.setNestedFileEnumeration(true) env.createInput(parquetInputFormat) .sinkTo(FileSink.forBulkFormat...) The job always fails with exception: java.nio.channels.ClosedChannelException at java.base/sun.nio.ch.FileChannelImpl.ensureOpen(FileChannelImpl.java:156) at java.base/sun.nio.ch.FileChannelImpl.position(FileChannelImpl.java:331) at org.apache.flink.core.fs.local.LocalRecoverableFsDataOutputStream.getPos(LocalRecoverableFsDataOutputStream.java:101) at org.apache.flink.formats.parquet.PositionOutputStreamAdapter.getPos(PositionOutputStreamAdapter.java:50) at org.apache.parquet.hadoop.ParquetFileWriter.startBlock(ParquetFileWriter.java:349) at org.apache.parquet.hadoop.InternalParquetRecordWriter.flushRowGroupToStore(InternalParquetRecordWriter.java:171) at org.apache.parquet.hadoop.InternalParquetRecordWriter.close(InternalParquetRecordWriter.java:114) at org.apache.parquet.hadoop.ParquetWriter.close(ParquetWriter.java:310) at org.apache.flink.formats.parquet.ParquetBulkWriter.finish(ParquetBulkWriter.java:62) at org.apache.flink.streaming.api.functions.sink.filesystem.BulkPartWriter.closeForCommit(BulkPartWriter.java:60) at org.apache.flink.connector.file.sink.writer.FileWriterBucket.closePartFile(FileWriterBucket.java:276) at org.apache.flink.connector.file.sink.writer.FileWriterBucket.prepareCommit(FileWriterBucket.java:202) at org.apache.flink.connector.file.sink.writer.FileWriter.prepareCommit(FileWriter.java:202) at org.apache.flink.streaming.runtime.operators.sink.AbstractSinkWriterOperator.endInput(AbstractSinkWriterOperator.java:97) at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.endOperatorInput(StreamOperatorWrapper.java:91) at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.lambda$close$0(StreamOperatorWrapper.java:127) at org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$1.runThrowing(StreamTaskActionExecutor.java:47) at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.close(StreamOperatorWrapper.java:127) at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.close(StreamOperatorWrapper.java:134) at org.apache.flink.streaming.runtime.tasks.OperatorChain.closeOperators(OperatorChain.java:412) at org.apache.flink.streaming.runtime.tasks.StreamTask.afterInvoke(StreamTask.java:585) at org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:547) at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:722) at org.apache.flink.runtime.taskmanager.Task.run(Task.java:547) This exception is always throwed after a warning: warn [ContinuousFileReaderOperator] not processing any records while closed I would supposed that problem is in my file sink, but the same file sink works for ParquetRowInputFormat. Did I miss something? Thank you! -- Sent from: http://apache-flink-user-mailing-list-archive.2336050.n4.nabble.com/ |
Update:
The job is working correctly if add an additional identity mapping step: env.createInput(parquetInputFormat) .map(record => record) .sinkTo(FileSink.forBulkFormat...) -- Sent from: http://apache-flink-user-mailing-list-archive.2336050.n4.nabble.com/ |
Hi Taras,
On first glance this looks like a bug to me. Can you try the latest 1.12 version (1.12.4)? If the bug still persists can you share the full job manager and task manager logs to further debug this problem. Best, Fabian > On 2. Jun 2021, at 13:22, Taras Moisiuk <[hidden email]> wrote: > > Update: > > The job is working correctly if add an additional identity mapping step: > > env.createInput(parquetInputFormat) > .map(record => record) > .sinkTo(FileSink.forBulkFormat...) > > > > -- > Sent from: http://apache-flink-user-mailing-list-archive.2336050.n4.nabble.com/ |
Hi Fabian,
Thank you for your answer. I've updated the flink version to 1.12.4 but unfortunately the problem still persists. I'm running this job in local mode, so I have only following log: WARNING: An illegal reflective access operation has occurred WARNING: Illegal reflective access by org.apache.flink.api.java.ClosureCleaner (file:/Users/user/Library/Caches/Coursier/v1/https/repo1.maven.org/maven2/org/apache/flink/flink-core/1.12.4/flink-core-1.12.4.jar) to field java.lang.String.value WARNING: Please consider reporting this to the maintainers of org.apache.flink.api.java.ClosureCleaner WARNING: Use --illegal-access=warn to enable warnings of further illegal reflective access operations WARNING: All illegal access operations will be denied in a future release 2021-06-02 17:13:00.713+0300 info [TypeExtractor] class org.apache.flink.streaming.api.functions.source.TimestampedFileInputSplit does not contain a setter for field modificationTime 2021-06-02 17:13:00.736+0300 info [TypeExtractor] Class class org.apache.flink.streaming.api.functions.source.TimestampedFileInputSplit cannot be used as a POJO type because not all fields are valid POJO fields, and must be processed as GenericType. Please read the Flink documentation on "Data Types & Serialization" for details of the effect on performance. 2021-06-02 17:13:01.136+0300 info [TaskExecutorResourceUtils] The configuration option taskmanager.cpu.cores required for local execution is not set, setting it to the maximal possible value. 2021-06-02 17:13:01.136+0300 info [TaskExecutorResourceUtils] The configuration option taskmanager.memory.task.heap.size required for local execution is not set, setting it to the maximal possible value. 2021-06-02 17:13:01.137+0300 info [TaskExecutorResourceUtils] The configuration option taskmanager.memory.task.off-heap.size required for local execution is not set, setting it to the maximal possible value. 2021-06-02 17:13:01.137+0300 info [TaskExecutorResourceUtils] The configuration option taskmanager.memory.network.min required for local execution is not set, setting it to its default value 64 mb. 2021-06-02 17:13:01.137+0300 info [TaskExecutorResourceUtils] The configuration option taskmanager.memory.network.max required for local execution is not set, setting it to its default value 64 mb. 2021-06-02 17:13:01.137+0300 info [TaskExecutorResourceUtils] The configuration option taskmanager.memory.managed.size required for local execution is not set, setting it to its default value 128 mb. 2021-06-02 17:13:01.147+0300 info [MiniCluster] Starting Flink Mini Cluster 2021-06-02 17:13:01.149+0300 info [MiniCluster] Starting Metrics Registry 2021-06-02 17:13:01.163+0300 info [MetricRegistryImpl] No metrics reporter configured, no metrics will be exposed/reported. 2021-06-02 17:13:01.164+0300 info [MiniCluster] Starting RPC Service(s) 2021-06-02 17:13:01.177+0300 info [AkkaRpcServiceUtils] Trying to start local actor system 2021-06-02 17:13:01.468+0300 info [Slf4jLogger] Slf4jLogger started 2021-06-02 17:13:01.621+0300 info [AkkaRpcServiceUtils] Actor system started at akka://flink 2021-06-02 17:13:01.638+0300 info [AkkaRpcServiceUtils] Trying to start local actor system 2021-06-02 17:13:01.651+0300 info [Slf4jLogger] Slf4jLogger started 2021-06-02 17:13:01.683+0300 info [AkkaRpcServiceUtils] Actor system started at akka://flink-metrics 2021-06-02 17:13:01.699+0300 info [AkkaRpcService] Starting RPC endpoint for org.apache.flink.runtime.metrics.dump.MetricQueryService at akka://flink-metrics/user/rpc/MetricQueryService . 2021-06-02 17:13:01.721+0300 info [MiniCluster] Starting high-availability services 2021-06-02 17:13:01.740+0300 info [BlobServer] Created BLOB server storage directory /var/folders/ms/18j3nxjd6rzcj9mmgly3n0rrq_g_vb/T/blobStore-880c269b-c4ef-475e-91a1-f542485ccee3 2021-06-02 17:13:01.749+0300 info [BlobServer] Started BLOB server at 0.0.0.0:56749 - max concurrent requests: 50 - max backlog: 1000 2021-06-02 17:13:01.756+0300 info [PermanentBlobCache] Created BLOB cache storage directory /var/folders/ms/18j3nxjd6rzcj9mmgly3n0rrq_g_vb/T/blobStore-1c57d8bd-0a31-464f-8c74-62ce32f56f49 2021-06-02 17:13:01.757+0300 info [TransientBlobCache] Created BLOB cache storage directory /var/folders/ms/18j3nxjd6rzcj9mmgly3n0rrq_g_vb/T/blobStore-6f8d5c1b-1ee5-49bd-bbba-c66aa8d18a90 2021-06-02 17:13:01.757+0300 info [MiniCluster] Starting 1 TaskManger(s) 2021-06-02 17:13:01.763+0300 info [TaskManagerRunner] Starting TaskManager with ResourceID: 264008fd-b7b7-4639-a3d4-ae430793106a 2021-06-02 17:13:01.784+0300 info [TaskManagerServices] Temporary file directory '/var/folders/ms/18j3nxjd6rzcj9mmgly3n0rrq_g_vb/T': total 233 GB, usable 18 GB (7.73% usable) 2021-06-02 17:13:01.789+0300 info [FileChannelManagerImpl] FileChannelManager uses directory /var/folders/ms/18j3nxjd6rzcj9mmgly3n0rrq_g_vb/T/flink-io-dcec76d9-4e46-411d-ad09-4257383a9df6 for spill files. 2021-06-02 17:13:01.800+0300 info [FileChannelManagerImpl] FileChannelManager uses directory /var/folders/ms/18j3nxjd6rzcj9mmgly3n0rrq_g_vb/T/flink-netty-shuffle-6412980d-faa6-4a3d-a3e0-138001f34bdc for spill files. 2021-06-02 17:13:01.851+0300 info [NetworkBufferPool] Allocated 64 MB for network buffer pool (number of memory segments: 2048, bytes per segment: 32768). 2021-06-02 17:13:01.866+0300 info [NettyShuffleEnvironment] Starting the network environment and its components. 2021-06-02 17:13:01.868+0300 info [KvStateService] Starting the kvState service and its components. 2021-06-02 17:13:01.898+0300 info [AkkaRpcService] Starting RPC endpoint for org.apache.flink.runtime.taskexecutor.TaskExecutor at akka://flink/user/rpc/taskmanager_0 . 2021-06-02 17:13:01.914+0300 info [DefaultJobLeaderService] Start job leader service. 2021-06-02 17:13:01.916+0300 info [FileCache] User file cache uses directory /var/folders/ms/18j3nxjd6rzcj9mmgly3n0rrq_g_vb/T/flink-dist-cache-763c1c92-6ae6-4e0d-8cec-2b42b11476ac 2021-06-02 17:13:01.972+0300 info [DispatcherRestEndpoint] Starting rest endpoint. 2021-06-02 17:13:02.273+0300 warn [WebMonitorUtils] Log file environment variable 'log.file' is not set. 2021-06-02 17:13:02.274+0300 warn [WebMonitorUtils] JobManager log files are unavailable in the web dashboard. Log file location not found in environment variable 'log.file' or configuration key 'web.log.path'. 2021-06-02 17:13:02.468+0300 info [DispatcherRestEndpoint] Rest endpoint listening at localhost:56750 2021-06-02 17:13:02.470+0300 info [EmbeddedLeaderService] Proposing leadership to contender http://localhost:56750 2021-06-02 17:13:02.472+0300 info [DispatcherRestEndpoint] Web frontend listening at http://localhost:56750. 2021-06-02 17:13:02.473+0300 info [DispatcherRestEndpoint] http://localhost:56750 was granted leadership with leaderSessionID=c200ef50-8c93-4a3f-b99e-ffe24618f62d 2021-06-02 17:13:02.473+0300 info [EmbeddedLeaderService] Received confirmation of leadership for leader http://localhost:56750 , session=c200ef50-8c93-4a3f-b99e-ffe24618f62d 2021-06-02 17:13:02.494+0300 info [AkkaRpcService] Starting RPC endpoint for org.apache.flink.runtime.resourcemanager.StandaloneResourceManager at akka://flink/user/rpc/resourcemanager_1 . 2021-06-02 17:13:02.514+0300 info [EmbeddedLeaderService] Proposing leadership to contender LeaderContender: DefaultDispatcherRunner 2021-06-02 17:13:02.515+0300 info [EmbeddedLeaderService] Proposing leadership to contender LeaderContender: StandaloneResourceManager 2021-06-02 17:13:02.516+0300 info [StandaloneResourceManager] ResourceManager akka://flink/user/rpc/resourcemanager_1 was granted leadership with fencing token 95c735454e4deaed3ef768b60cb64d55 2021-06-02 17:13:02.519+0300 info [MiniCluster] Flink Mini Cluster started successfully 2021-06-02 17:13:02.520+0300 info [SlotManagerImpl] Starting the SlotManager. 2021-06-02 17:13:02.520+0300 info [SessionDispatcherLeaderProcess] Start SessionDispatcherLeaderProcess. 2021-06-02 17:13:02.521+0300 info [SessionDispatcherLeaderProcess] Recover all persisted job graphs. 2021-06-02 17:13:02.522+0300 info [SessionDispatcherLeaderProcess] Successfully recovered 0 persisted job graphs. 2021-06-02 17:13:02.528+0300 info [EmbeddedLeaderService] Received confirmation of leadership for leader akka://flink/user/rpc/resourcemanager_1 , session=3ef768b6-0cb6-4d55-95c7-35454e4deaed 2021-06-02 17:13:02.529+0300 info [TaskExecutor] Connecting to ResourceManager akka://flink/user/rpc/resourcemanager_1(95c735454e4deaed3ef768b60cb64d55). 2021-06-02 17:13:02.533+0300 info [AkkaRpcService] Starting RPC endpoint for org.apache.flink.runtime.dispatcher.StandaloneDispatcher at akka://flink/user/rpc/dispatcher_2 . 2021-06-02 17:13:02.544+0300 info [EmbeddedLeaderService] Received confirmation of leadership for leader akka://flink/user/rpc/dispatcher_2 , session=504bd90b-fd4b-43e9-bc71-ac0d30fe13a3 2021-06-02 17:13:02.551+0300 info [TaskExecutor] Resolved ResourceManager address, beginning registration 2021-06-02 17:13:02.559+0300 info [StandaloneResourceManager] Registering TaskManager with ResourceID 264008fd-b7b7-4639-a3d4-ae430793106a (akka://flink/user/rpc/taskmanager_0) at ResourceManager 2021-06-02 17:13:02.561+0300 info [TaskExecutor] Successful registration at resource manager akka://flink/user/rpc/resourcemanager_1 under registration id 0f21d7e84e09c3d6ad31bfa99396c744. 2021-06-02 17:13:02.562+0300 info [StandaloneDispatcher] Received JobGraph submission d70f7e1ff824fb485de8a4dbf2b485b0 (hdfs_parquet_compacter). 2021-06-02 17:13:02.563+0300 info [StandaloneDispatcher] Submitting job d70f7e1ff824fb485de8a4dbf2b485b0 (hdfs_parquet_compacter). 2021-06-02 17:13:02.585+0300 info [EmbeddedLeaderService] Proposing leadership to contender LeaderContender: JobManagerRunnerImpl 2021-06-02 17:13:02.596+0300 info [AkkaRpcService] Starting RPC endpoint for org.apache.flink.runtime.jobmaster.JobMaster at akka://flink/user/rpc/jobmanager_3 . 2021-06-02 17:13:02.606+0300 info [JobMaster] Initializing job hdfs_parquet_compacter (d70f7e1ff824fb485de8a4dbf2b485b0). 2021-06-02 17:13:02.637+0300 info [JobMaster] Using restart back off time strategy NoRestartBackoffTimeStrategy for hdfs_parquet_compacter (d70f7e1ff824fb485de8a4dbf2b485b0). 2021-06-02 17:13:02.686+0300 info [JobMaster] Running initialization on master for job hdfs_parquet_compacter (d70f7e1ff824fb485de8a4dbf2b485b0). 2021-06-02 17:13:02.687+0300 info [JobMaster] Successfully ran initialization on master in 0 ms. 2021-06-02 17:13:02.704+0300 info [DefaultExecutionTopology] Built 10 pipelined regions in 0 ms 2021-06-02 17:13:02.718+0300 info [JobMaster] Using application-defined state backend: org.apache.flink.streaming.api.operators.sorted.state.BatchExecutionStateBackend@77483711 2021-06-02 17:13:02.730+0300 info [CheckpointCoordinator] No checkpoint found during restore. 2021-06-02 17:13:02.734+0300 info [JobMaster] Using failover strategy org.apache.flink.runtime.executiongraph.failover.flip1.RestartPipelinedRegionFailoverStrategy@4627cc71 for hdfs_parquet_compacter (d70f7e1ff824fb485de8a4dbf2b485b0). 2021-06-02 17:13:02.745+0300 info [JobManagerRunnerImpl] JobManager runner for job hdfs_parquet_compacter (d70f7e1ff824fb485de8a4dbf2b485b0) was granted leadership with session id 55764aae-a3b9-435d-bc02-4d08781064bc at akka://flink/user/rpc/jobmanager_3. 2021-06-02 17:13:02.748+0300 info [JobMaster] Starting execution of job hdfs_parquet_compacter (d70f7e1ff824fb485de8a4dbf2b485b0) under job master id bc024d08781064bc55764aaea3b9435d. 2021-06-02 17:13:02.749+0300 info [JobMaster] Starting scheduling with scheduling strategy [org.apache.flink.runtime.scheduler.strategy.PipelinedRegionSchedulingStrategy] 2021-06-02 17:13:02.750+0300 info [ExecutionGraph] Job hdfs_parquet_compacter (d70f7e1ff824fb485de8a4dbf2b485b0) switched from state CREATED to RUNNING. 2021-06-02 17:13:02.754+0300 info [ExecutionGraph] Source: Custom File source (1/1) (6330e22e7399f04ba7087e1b2e885044) switched from CREATED to SCHEDULED. 2021-06-02 17:13:02.767+0300 info [SlotPoolImpl] Cannot serve slot request, no ResourceManager connected. Adding as pending request [SlotRequestId{6b21c901901dc676dfcf757376901af8}] 2021-06-02 17:13:02.774+0300 info [EmbeddedLeaderService] Received confirmation of leadership for leader akka://flink/user/rpc/jobmanager_3 , session=55764aae-a3b9-435d-bc02-4d08781064bc 2021-06-02 17:13:02.774+0300 info [JobMaster] Connecting to ResourceManager akka://flink/user/rpc/resourcemanager_1(95c735454e4deaed3ef768b60cb64d55) 2021-06-02 17:13:02.776+0300 info [JobMaster] Resolved ResourceManager address, beginning registration 2021-06-02 17:13:02.778+0300 info [StandaloneResourceManager] Registering job manager bc024d08781064bc55764aaea3b9435d@akka://flink/user/rpc/jobmanager_3 for job d70f7e1ff824fb485de8a4dbf2b485b0. 2021-06-02 17:13:02.783+0300 info [StandaloneResourceManager] Registered job manager bc024d08781064bc55764aaea3b9435d@akka://flink/user/rpc/jobmanager_3 for job d70f7e1ff824fb485de8a4dbf2b485b0. 2021-06-02 17:13:02.785+0300 info [JobMaster] JobManager successfully registered at ResourceManager, leader id: 95c735454e4deaed3ef768b60cb64d55. 2021-06-02 17:13:02.786+0300 info [SlotPoolImpl] Requesting new slot [SlotRequestId{6b21c901901dc676dfcf757376901af8}] and profile ResourceProfile{UNKNOWN} with allocation id 2b3f94c62b0602f21dd1b4af99bec0fc from resource manager. 2021-06-02 17:13:02.787+0300 info [StandaloneResourceManager] Request slot with profile ResourceProfile{UNKNOWN} for job d70f7e1ff824fb485de8a4dbf2b485b0 with allocation id 2b3f94c62b0602f21dd1b4af99bec0fc. 2021-06-02 17:13:02.789+0300 info [TaskExecutor] Receive slot request 2b3f94c62b0602f21dd1b4af99bec0fc for job d70f7e1ff824fb485de8a4dbf2b485b0 from resource manager with leader id 95c735454e4deaed3ef768b60cb64d55. 2021-06-02 17:13:02.796+0300 info [TaskExecutor] Allocated slot for 2b3f94c62b0602f21dd1b4af99bec0fc. 2021-06-02 17:13:02.797+0300 info [DefaultJobLeaderService] Add job d70f7e1ff824fb485de8a4dbf2b485b0 for job leader monitoring. 2021-06-02 17:13:02.799+0300 info [DefaultJobLeaderService] Try to register at job manager akka://flink/user/rpc/jobmanager_3 with leader id 55764aae-a3b9-435d-bc02-4d08781064bc. 2021-06-02 17:13:02.799+0300 info [DefaultJobLeaderService] Resolved JobManager address, beginning registration 2021-06-02 17:13:02.804+0300 info [DefaultJobLeaderService] Successful registration at job manager akka://flink/user/rpc/jobmanager_3 for job d70f7e1ff824fb485de8a4dbf2b485b0. 2021-06-02 17:13:02.805+0300 info [TaskExecutor] Establish JobManager connection for job d70f7e1ff824fb485de8a4dbf2b485b0. 2021-06-02 17:13:02.809+0300 info [TaskExecutor] Offer reserved slots to the leader of job d70f7e1ff824fb485de8a4dbf2b485b0. 2021-06-02 17:13:02.817+0300 info [ExecutionGraph] Source: Custom File source (1/1) (6330e22e7399f04ba7087e1b2e885044) switched from SCHEDULED to DEPLOYING. 2021-06-02 17:13:02.818+0300 info [ExecutionGraph] Deploying Source: Custom File source (1/1) (attempt #0) with attempt id 6330e22e7399f04ba7087e1b2e885044 to 264008fd-b7b7-4639-a3d4-ae430793106a @ localhost (dataPort=-1) with allocation id 2b3f94c62b0602f21dd1b4af99bec0fc 2021-06-02 17:13:02.822+0300 info [TaskSlotTableImpl] Activate slot 2b3f94c62b0602f21dd1b4af99bec0fc. 2021-06-02 17:13:02.853+0300 info [TaskExecutor] Received task Source: Custom File source (1/1)#0 (6330e22e7399f04ba7087e1b2e885044), deploy into slot with allocation id 2b3f94c62b0602f21dd1b4af99bec0fc. 2021-06-02 17:13:02.854+0300 info [Task] Source: Custom File source (1/1)#0 (6330e22e7399f04ba7087e1b2e885044) switched from CREATED to DEPLOYING. 2021-06-02 17:13:02.856+0300 info [Task] Loading JAR files for task Source: Custom File source (1/1)#0 (6330e22e7399f04ba7087e1b2e885044) [DEPLOYING]. 2021-06-02 17:13:02.857+0300 info [TaskSlotTableImpl] Activate slot 2b3f94c62b0602f21dd1b4af99bec0fc. 2021-06-02 17:13:02.858+0300 info [Task] Registering task at network: Source: Custom File source (1/1)#0 (6330e22e7399f04ba7087e1b2e885044) [DEPLOYING]. 2021-06-02 17:13:02.874+0300 info [StreamTask] Using application-defined state backend: org.apache.flink.streaming.api.operators.sorted.state.BatchExecutionStateBackend@17b60aba 2021-06-02 17:13:02.881+0300 info [Task] Source: Custom File source (1/1)#0 (6330e22e7399f04ba7087e1b2e885044) switched from DEPLOYING to RUNNING. 2021-06-02 17:13:02.882+0300 info [ExecutionGraph] Source: Custom File source (1/1) (6330e22e7399f04ba7087e1b2e885044) switched from DEPLOYING to RUNNING. 2021-06-02 17:13:02.937+0300 info [ContinuousFileMonitoringFunction] No state to restore for the ContinuousFileMonitoringFunction. 2021-06-02 17:13:02.946+0300 info [ContinuousFileMonitoringFunction] Forwarding split: [0] file:/Users/user/compacter/tmp/parquet_example/kafka_timestamp=2021-05-19/part-efbb9d61-1319-48e0-8682-1c5701e7bf1d-0 mod@ 1622197611838 : 0 + 18766584 2021-06-02 17:13:03.113+0300 info [Task] Source: Custom File source (1/1)#0 (6330e22e7399f04ba7087e1b2e885044) switched from RUNNING to FINISHED. 2021-06-02 17:13:03.113+0300 info [Task] Freeing task resources for Source: Custom File source (1/1)#0 (6330e22e7399f04ba7087e1b2e885044). 2021-06-02 17:13:03.114+0300 info [TaskExecutor] Un-registering task and sending final execution state FINISHED to JobManager for task Source: Custom File source (1/1)#0 6330e22e7399f04ba7087e1b2e885044. 2021-06-02 17:13:03.120+0300 info [ExecutionGraph] Source: Custom File source (1/1) (6330e22e7399f04ba7087e1b2e885044) switched from RUNNING to FINISHED. 2021-06-02 17:13:03.125+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (1/8) (198c68d2f177dda7d3f1803a5c8283ca) switched from CREATED to SCHEDULED. 2021-06-02 17:13:03.127+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (1/8) (198c68d2f177dda7d3f1803a5c8283ca) switched from SCHEDULED to DEPLOYING. 2021-06-02 17:13:03.127+0300 info [ExecutionGraph] Deploying /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (1/8) (attempt #0) with attempt id 198c68d2f177dda7d3f1803a5c8283ca to 264008fd-b7b7-4639-a3d4-ae430793106a @ localhost (dataPort=-1) with allocation id 2b3f94c62b0602f21dd1b4af99bec0fc 2021-06-02 17:13:03.129+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (2/8) (143e6297eb2218f5590cf068e386ba71) switched from CREATED to SCHEDULED. 2021-06-02 17:13:03.129+0300 info [TaskSlotTableImpl] Activate slot 2b3f94c62b0602f21dd1b4af99bec0fc. 2021-06-02 17:13:03.129+0300 info [SlotPoolImpl] Requesting new slot [SlotRequestId{2a1e79f4e25924f707c9dcd74e8c50d3}] and profile ResourceProfile{UNKNOWN} with allocation id 5a88e10c6d023fed666e644c9af0a210 from resource manager. 2021-06-02 17:13:03.130+0300 info [StandaloneResourceManager] Request slot with profile ResourceProfile{UNKNOWN} for job d70f7e1ff824fb485de8a4dbf2b485b0 with allocation id 5a88e10c6d023fed666e644c9af0a210. 2021-06-02 17:13:03.130+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (3/8) (2906d5a3cbb74007054260ac3edf6385) switched from CREATED to SCHEDULED. 2021-06-02 17:13:03.130+0300 info [SlotPoolImpl] Requesting new slot [SlotRequestId{b7101290c5a88a4d0ac68f81ceaf016b}] and profile ResourceProfile{UNKNOWN} with allocation id 96738ae5a01c330b7899819edb6918b6 from resource manager. 2021-06-02 17:13:03.130+0300 info [StandaloneResourceManager] Request slot with profile ResourceProfile{UNKNOWN} for job d70f7e1ff824fb485de8a4dbf2b485b0 with allocation id 96738ae5a01c330b7899819edb6918b6. 2021-06-02 17:13:03.131+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (4/8) (1c884a8478c56d4a9270d30c140724e8) switched from CREATED to SCHEDULED. 2021-06-02 17:13:03.131+0300 info [SlotPoolImpl] Requesting new slot [SlotRequestId{2365594bb95672c4f08c1ca80fc19ba9}] and profile ResourceProfile{UNKNOWN} with allocation id 43d084de5b721547c53d65b821d01f6f from resource manager. 2021-06-02 17:13:03.131+0300 info [StandaloneResourceManager] Request slot with profile ResourceProfile{UNKNOWN} for job d70f7e1ff824fb485de8a4dbf2b485b0 with allocation id 43d084de5b721547c53d65b821d01f6f. 2021-06-02 17:13:03.132+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (5/8) (d8fa853f251afd4f87776ad0b3810635) switched from CREATED to SCHEDULED. 2021-06-02 17:13:03.132+0300 info [SlotPoolImpl] Requesting new slot [SlotRequestId{d28fbd227b1770dc533162c55075b53e}] and profile ResourceProfile{UNKNOWN} with allocation id 0db08c17830fa9b358fd76153487f39d from resource manager. 2021-06-02 17:13:03.133+0300 info [StandaloneResourceManager] Request slot with profile ResourceProfile{UNKNOWN} for job d70f7e1ff824fb485de8a4dbf2b485b0 with allocation id 0db08c17830fa9b358fd76153487f39d. 2021-06-02 17:13:03.133+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (6/8) (7e273ba7189fa887157b7cf00565aa27) switched from CREATED to SCHEDULED. 2021-06-02 17:13:03.133+0300 info [SlotPoolImpl] Requesting new slot [SlotRequestId{74b87314e91f1ed064e1ad035d7c6a8f}] and profile ResourceProfile{UNKNOWN} with allocation id 237b76cbdd6a12a219a9600ab9fb667a from resource manager. 2021-06-02 17:13:03.134+0300 info [StandaloneResourceManager] Request slot with profile ResourceProfile{UNKNOWN} for job d70f7e1ff824fb485de8a4dbf2b485b0 with allocation id 237b76cbdd6a12a219a9600ab9fb667a. 2021-06-02 17:13:03.134+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (7/8) (9b4b5cf10ffef54e899ab3c200403a3c) switched from CREATED to SCHEDULED. 2021-06-02 17:13:03.134+0300 info [SlotPoolImpl] Requesting new slot [SlotRequestId{313fdb262f6885e7668ef68a24c1b4e4}] and profile ResourceProfile{UNKNOWN} with allocation id d12f6d0d0893e1b326e2b5e94818ea06 from resource manager. 2021-06-02 17:13:03.135+0300 info [StandaloneResourceManager] Request slot with profile ResourceProfile{UNKNOWN} for job d70f7e1ff824fb485de8a4dbf2b485b0 with allocation id d12f6d0d0893e1b326e2b5e94818ea06. 2021-06-02 17:13:03.135+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (8/8) (ba0cbe884bf9462ac3ef39f7455be895) switched from CREATED to SCHEDULED. 2021-06-02 17:13:03.135+0300 info [SlotPoolImpl] Requesting new slot [SlotRequestId{824c94987f1fa3f712eb8f22936ba8d5}] and profile ResourceProfile{UNKNOWN} with allocation id 8ab384bad8ddea1d6796c14ec8ffca36 from resource manager. 2021-06-02 17:13:03.136+0300 info [StandaloneResourceManager] Request slot with profile ResourceProfile{UNKNOWN} for job d70f7e1ff824fb485de8a4dbf2b485b0 with allocation id 8ab384bad8ddea1d6796c14ec8ffca36. 2021-06-02 17:13:03.142+0300 info [TaskExecutor] Received task /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (1/8)#0 (198c68d2f177dda7d3f1803a5c8283ca), deploy into slot with allocation id 2b3f94c62b0602f21dd1b4af99bec0fc. 2021-06-02 17:13:03.142+0300 info [Task] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (1/8)#0 (198c68d2f177dda7d3f1803a5c8283ca) switched from CREATED to DEPLOYING. 2021-06-02 17:13:03.143+0300 info [TaskExecutor] Receive slot request 5a88e10c6d023fed666e644c9af0a210 for job d70f7e1ff824fb485de8a4dbf2b485b0 from resource manager with leader id 95c735454e4deaed3ef768b60cb64d55. 2021-06-02 17:13:03.143+0300 info [Task] Loading JAR files for task /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (1/8)#0 (198c68d2f177dda7d3f1803a5c8283ca) [DEPLOYING]. 2021-06-02 17:13:03.143+0300 info [TaskExecutor] Allocated slot for 5a88e10c6d023fed666e644c9af0a210. 2021-06-02 17:13:03.144+0300 info [TaskExecutor] Offer reserved slots to the leader of job d70f7e1ff824fb485de8a4dbf2b485b0. 2021-06-02 17:13:03.144+0300 info [Task] Registering task at network: /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (1/8)#0 (198c68d2f177dda7d3f1803a5c8283ca) [DEPLOYING]. 2021-06-02 17:13:03.144+0300 info [TaskExecutor] Receive slot request 96738ae5a01c330b7899819edb6918b6 for job d70f7e1ff824fb485de8a4dbf2b485b0 from resource manager with leader id 95c735454e4deaed3ef768b60cb64d55. 2021-06-02 17:13:03.145+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (2/8) (143e6297eb2218f5590cf068e386ba71) switched from SCHEDULED to DEPLOYING. 2021-06-02 17:13:03.145+0300 info [TaskExecutor] Allocated slot for 96738ae5a01c330b7899819edb6918b6. 2021-06-02 17:13:03.146+0300 info [ExecutionGraph] Deploying /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (2/8) (attempt #0) with attempt id 143e6297eb2218f5590cf068e386ba71 to 264008fd-b7b7-4639-a3d4-ae430793106a @ localhost (dataPort=-1) with allocation id 5a88e10c6d023fed666e644c9af0a210 2021-06-02 17:13:03.146+0300 info [TaskExecutor] Offer reserved slots to the leader of job d70f7e1ff824fb485de8a4dbf2b485b0. 2021-06-02 17:13:03.146+0300 info [TaskExecutor] Receive slot request 43d084de5b721547c53d65b821d01f6f for job d70f7e1ff824fb485de8a4dbf2b485b0 from resource manager with leader id 95c735454e4deaed3ef768b60cb64d55. 2021-06-02 17:13:03.146+0300 info [TaskExecutor] Allocated slot for 43d084de5b721547c53d65b821d01f6f. 2021-06-02 17:13:03.147+0300 info [TaskExecutor] Offer reserved slots to the leader of job d70f7e1ff824fb485de8a4dbf2b485b0. 2021-06-02 17:13:03.147+0300 info [SlotPoolImpl] Received repeated offer for slot [5a88e10c6d023fed666e644c9af0a210]. Ignoring. 2021-06-02 17:13:03.147+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (3/8) (2906d5a3cbb74007054260ac3edf6385) switched from SCHEDULED to DEPLOYING. 2021-06-02 17:13:03.147+0300 info [ExecutionGraph] Deploying /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (3/8) (attempt #0) with attempt id 2906d5a3cbb74007054260ac3edf6385 to 264008fd-b7b7-4639-a3d4-ae430793106a @ localhost (dataPort=-1) with allocation id 96738ae5a01c330b7899819edb6918b6 2021-06-02 17:13:03.148+0300 info [TaskExecutor] Receive slot request 0db08c17830fa9b358fd76153487f39d for job d70f7e1ff824fb485de8a4dbf2b485b0 from resource manager with leader id 95c735454e4deaed3ef768b60cb64d55. 2021-06-02 17:13:03.148+0300 info [TaskExecutor] Allocated slot for 0db08c17830fa9b358fd76153487f39d. 2021-06-02 17:13:03.149+0300 info [TaskExecutor] Offer reserved slots to the leader of job d70f7e1ff824fb485de8a4dbf2b485b0. 2021-06-02 17:13:03.149+0300 info [StreamTask] Using application-defined state backend: org.apache.flink.streaming.api.operators.sorted.state.BatchExecutionStateBackend@3d7cbf83 2021-06-02 17:13:03.149+0300 info [Task] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (1/8)#0 (198c68d2f177dda7d3f1803a5c8283ca) switched from DEPLOYING to RUNNING. 2021-06-02 17:13:03.149+0300 info [TaskExecutor] Receive slot request 237b76cbdd6a12a219a9600ab9fb667a for job d70f7e1ff824fb485de8a4dbf2b485b0 from resource manager with leader id 95c735454e4deaed3ef768b60cb64d55. 2021-06-02 17:13:03.150+0300 info [SlotPoolImpl] Received repeated offer for slot [5a88e10c6d023fed666e644c9af0a210]. Ignoring. 2021-06-02 17:13:03.150+0300 info [TaskExecutor] Allocated slot for 237b76cbdd6a12a219a9600ab9fb667a. 2021-06-02 17:13:03.150+0300 info [SlotPoolImpl] Received repeated offer for slot [96738ae5a01c330b7899819edb6918b6]. Ignoring. 2021-06-02 17:13:03.150+0300 info [TaskExecutor] Offer reserved slots to the leader of job d70f7e1ff824fb485de8a4dbf2b485b0. 2021-06-02 17:13:03.151+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (4/8) (1c884a8478c56d4a9270d30c140724e8) switched from SCHEDULED to DEPLOYING. 2021-06-02 17:13:03.151+0300 info [ExecutionGraph] Deploying /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (4/8) (attempt #0) with attempt id 1c884a8478c56d4a9270d30c140724e8 to 264008fd-b7b7-4639-a3d4-ae430793106a @ localhost (dataPort=-1) with allocation id 43d084de5b721547c53d65b821d01f6f 2021-06-02 17:13:03.151+0300 info [TaskExecutor] Receive slot request d12f6d0d0893e1b326e2b5e94818ea06 for job d70f7e1ff824fb485de8a4dbf2b485b0 from resource manager with leader id 95c735454e4deaed3ef768b60cb64d55. 2021-06-02 17:13:03.151+0300 info [TaskExecutor] Allocated slot for d12f6d0d0893e1b326e2b5e94818ea06. 2021-06-02 17:13:03.151+0300 info [TaskExecutor] Offer reserved slots to the leader of job d70f7e1ff824fb485de8a4dbf2b485b0. 2021-06-02 17:13:03.152+0300 info [TaskExecutor] Receive slot request 8ab384bad8ddea1d6796c14ec8ffca36 for job d70f7e1ff824fb485de8a4dbf2b485b0 from resource manager with leader id 95c735454e4deaed3ef768b60cb64d55. 2021-06-02 17:13:03.152+0300 info [TaskExecutor] Allocated slot for 8ab384bad8ddea1d6796c14ec8ffca36. 2021-06-02 17:13:03.152+0300 info [TaskExecutor] Offer reserved slots to the leader of job d70f7e1ff824fb485de8a4dbf2b485b0. 2021-06-02 17:13:03.152+0300 info [TaskSlotTableImpl] Activate slot 5a88e10c6d023fed666e644c9af0a210. 2021-06-02 17:13:03.152+0300 info [SlotPoolImpl] Received repeated offer for slot [5a88e10c6d023fed666e644c9af0a210]. Ignoring. 2021-06-02 17:13:03.153+0300 info [SlotPoolImpl] Received repeated offer for slot [96738ae5a01c330b7899819edb6918b6]. Ignoring. 2021-06-02 17:13:03.153+0300 info [SlotPoolImpl] Received repeated offer for slot [43d084de5b721547c53d65b821d01f6f]. Ignoring. 2021-06-02 17:13:03.153+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (5/8) (d8fa853f251afd4f87776ad0b3810635) switched from SCHEDULED to DEPLOYING. 2021-06-02 17:13:03.153+0300 info [ExecutionGraph] Deploying /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (5/8) (attempt #0) with attempt id d8fa853f251afd4f87776ad0b3810635 to 264008fd-b7b7-4639-a3d4-ae430793106a @ localhost (dataPort=-1) with allocation id 0db08c17830fa9b358fd76153487f39d 2021-06-02 17:13:03.153+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (1/8) (198c68d2f177dda7d3f1803a5c8283ca) switched from DEPLOYING to RUNNING. 2021-06-02 17:13:03.154+0300 info [SlotPoolImpl] Received repeated offer for slot [5a88e10c6d023fed666e644c9af0a210]. Ignoring. 2021-06-02 17:13:03.154+0300 info [SlotPoolImpl] Received repeated offer for slot [96738ae5a01c330b7899819edb6918b6]. Ignoring. 2021-06-02 17:13:03.154+0300 info [SlotPoolImpl] Received repeated offer for slot [43d084de5b721547c53d65b821d01f6f]. Ignoring. 2021-06-02 17:13:03.154+0300 info [SlotPoolImpl] Received repeated offer for slot [0db08c17830fa9b358fd76153487f39d]. Ignoring. 2021-06-02 17:13:03.154+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (6/8) (7e273ba7189fa887157b7cf00565aa27) switched from SCHEDULED to DEPLOYING. 2021-06-02 17:13:03.155+0300 info [ExecutionGraph] Deploying /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (6/8) (attempt #0) with attempt id 7e273ba7189fa887157b7cf00565aa27 to 264008fd-b7b7-4639-a3d4-ae430793106a @ localhost (dataPort=-1) with allocation id 237b76cbdd6a12a219a9600ab9fb667a 2021-06-02 17:13:03.155+0300 info [TaskExecutor] Received task /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (2/8)#0 (143e6297eb2218f5590cf068e386ba71), deploy into slot with allocation id 5a88e10c6d023fed666e644c9af0a210. 2021-06-02 17:13:03.155+0300 info [SlotPoolImpl] Received repeated offer for slot [43d084de5b721547c53d65b821d01f6f]. Ignoring. 2021-06-02 17:13:03.155+0300 info [SlotPoolImpl] Received repeated offer for slot [5a88e10c6d023fed666e644c9af0a210]. Ignoring. 2021-06-02 17:13:03.156+0300 info [SlotPoolImpl] Received repeated offer for slot [96738ae5a01c330b7899819edb6918b6]. Ignoring. 2021-06-02 17:13:03.156+0300 info [SlotPoolImpl] Received repeated offer for slot [0db08c17830fa9b358fd76153487f39d]. Ignoring. 2021-06-02 17:13:03.156+0300 info [SlotPoolImpl] Received repeated offer for slot [237b76cbdd6a12a219a9600ab9fb667a]. Ignoring. 2021-06-02 17:13:03.156+0300 info [Task] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (2/8)#0 (143e6297eb2218f5590cf068e386ba71) switched from CREATED to DEPLOYING. 2021-06-02 17:13:03.156+0300 info [TaskSlotTableImpl] Activate slot 5a88e10c6d023fed666e644c9af0a210. 2021-06-02 17:13:03.156+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (7/8) (9b4b5cf10ffef54e899ab3c200403a3c) switched from SCHEDULED to DEPLOYING. 2021-06-02 17:13:03.157+0300 info [ExecutionGraph] Deploying /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (7/8) (attempt #0) with attempt id 9b4b5cf10ffef54e899ab3c200403a3c to 264008fd-b7b7-4639-a3d4-ae430793106a @ localhost (dataPort=-1) with allocation id d12f6d0d0893e1b326e2b5e94818ea06 2021-06-02 17:13:03.157+0300 info [TaskSlotTableImpl] Activate slot 96738ae5a01c330b7899819edb6918b6. 2021-06-02 17:13:03.157+0300 info [Task] Loading JAR files for task /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (2/8)#0 (143e6297eb2218f5590cf068e386ba71) [DEPLOYING]. 2021-06-02 17:13:03.157+0300 info [SlotPoolImpl] Received repeated offer for slot [43d084de5b721547c53d65b821d01f6f]. Ignoring. 2021-06-02 17:13:03.157+0300 info [SlotPoolImpl] Received repeated offer for slot [d12f6d0d0893e1b326e2b5e94818ea06]. Ignoring. 2021-06-02 17:13:03.158+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (8/8) (ba0cbe884bf9462ac3ef39f7455be895) switched from SCHEDULED to DEPLOYING. 2021-06-02 17:13:03.158+0300 info [ExecutionGraph] Deploying /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (8/8) (attempt #0) with attempt id ba0cbe884bf9462ac3ef39f7455be895 to 264008fd-b7b7-4639-a3d4-ae430793106a @ localhost (dataPort=-1) with allocation id 8ab384bad8ddea1d6796c14ec8ffca36 2021-06-02 17:13:03.158+0300 info [Task] Registering task at network: /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (2/8)#0 (143e6297eb2218f5590cf068e386ba71) [DEPLOYING]. 2021-06-02 17:13:03.158+0300 info [SlotPoolImpl] Received repeated offer for slot [5a88e10c6d023fed666e644c9af0a210]. Ignoring. 2021-06-02 17:13:03.158+0300 info [SlotPoolImpl] Received repeated offer for slot [96738ae5a01c330b7899819edb6918b6]. Ignoring. 2021-06-02 17:13:03.158+0300 info [SlotPoolImpl] Received repeated offer for slot [0db08c17830fa9b358fd76153487f39d]. Ignoring. 2021-06-02 17:13:03.159+0300 info [SlotPoolImpl] Received repeated offer for slot [237b76cbdd6a12a219a9600ab9fb667a]. Ignoring. 2021-06-02 17:13:03.159+0300 info [StreamTask] Using application-defined state backend: org.apache.flink.streaming.api.operators.sorted.state.BatchExecutionStateBackend@424e6bfc 2021-06-02 17:13:03.159+0300 info [Task] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (2/8)#0 (143e6297eb2218f5590cf068e386ba71) switched from DEPLOYING to RUNNING. 2021-06-02 17:13:03.159+0300 info [TaskExecutor] Received task /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (3/8)#0 (2906d5a3cbb74007054260ac3edf6385), deploy into slot with allocation id 96738ae5a01c330b7899819edb6918b6. 2021-06-02 17:13:03.159+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (2/8) (143e6297eb2218f5590cf068e386ba71) switched from DEPLOYING to RUNNING. 2021-06-02 17:13:03.160+0300 info [Task] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (3/8)#0 (2906d5a3cbb74007054260ac3edf6385) switched from CREATED to DEPLOYING. 2021-06-02 17:13:03.160+0300 info [Task] Loading JAR files for task /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (3/8)#0 (2906d5a3cbb74007054260ac3edf6385) [DEPLOYING]. 2021-06-02 17:13:03.160+0300 info [TaskSlotTableImpl] Activate slot 5a88e10c6d023fed666e644c9af0a210. 2021-06-02 17:13:03.160+0300 info [TaskSlotTableImpl] Activate slot 96738ae5a01c330b7899819edb6918b6. 2021-06-02 17:13:03.161+0300 info [TaskSlotTableImpl] Activate slot 43d084de5b721547c53d65b821d01f6f. 2021-06-02 17:13:03.162+0300 info [Task] Registering task at network: /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (3/8)#0 (2906d5a3cbb74007054260ac3edf6385) [DEPLOYING]. 2021-06-02 17:13:03.163+0300 info [StreamTask] Using application-defined state backend: org.apache.flink.streaming.api.operators.sorted.state.BatchExecutionStateBackend@72f44e2f 2021-06-02 17:13:03.164+0300 info [Task] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (3/8)#0 (2906d5a3cbb74007054260ac3edf6385) switched from DEPLOYING to RUNNING. 2021-06-02 17:13:03.164+0300 info [TaskExecutor] Received task /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (4/8)#0 (1c884a8478c56d4a9270d30c140724e8), deploy into slot with allocation id 43d084de5b721547c53d65b821d01f6f. 2021-06-02 17:13:03.165+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (3/8) (2906d5a3cbb74007054260ac3edf6385) switched from DEPLOYING to RUNNING. 2021-06-02 17:13:03.166+0300 info [Task] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (4/8)#0 (1c884a8478c56d4a9270d30c140724e8) switched from CREATED to DEPLOYING. 2021-06-02 17:13:03.166+0300 info [Task] Loading JAR files for task /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (4/8)#0 (1c884a8478c56d4a9270d30c140724e8) [DEPLOYING]. 2021-06-02 17:13:03.166+0300 info [Task] Registering task at network: /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (4/8)#0 (1c884a8478c56d4a9270d30c140724e8) [DEPLOYING]. 2021-06-02 17:13:03.169+0300 info [StreamTask] Using application-defined state backend: org.apache.flink.streaming.api.operators.sorted.state.BatchExecutionStateBackend@2e50f4be 2021-06-02 17:13:03.170+0300 info [Task] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (4/8)#0 (1c884a8478c56d4a9270d30c140724e8) switched from DEPLOYING to RUNNING. 2021-06-02 17:13:03.170+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (4/8) (1c884a8478c56d4a9270d30c140724e8) switched from DEPLOYING to RUNNING. 2021-06-02 17:13:03.170+0300 info [TaskSlotTableImpl] Activate slot 5a88e10c6d023fed666e644c9af0a210. 2021-06-02 17:13:03.170+0300 info [TaskSlotTableImpl] Activate slot 96738ae5a01c330b7899819edb6918b6. 2021-06-02 17:13:03.171+0300 info [TaskSlotTableImpl] Activate slot 43d084de5b721547c53d65b821d01f6f. 2021-06-02 17:13:03.171+0300 info [TaskSlotTableImpl] Activate slot 5a88e10c6d023fed666e644c9af0a210. 2021-06-02 17:13:03.171+0300 info [TaskSlotTableImpl] Activate slot 96738ae5a01c330b7899819edb6918b6. 2021-06-02 17:13:03.171+0300 info [TaskSlotTableImpl] Activate slot 43d084de5b721547c53d65b821d01f6f. 2021-06-02 17:13:03.171+0300 info [TaskSlotTableImpl] Activate slot 0db08c17830fa9b358fd76153487f39d. 2021-06-02 17:13:03.172+0300 info [TaskSlotTableImpl] Activate slot 0db08c17830fa9b358fd76153487f39d. 2021-06-02 17:13:03.173+0300 info [TaskExecutor] Received task /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (5/8)#0 (d8fa853f251afd4f87776ad0b3810635), deploy into slot with allocation id 0db08c17830fa9b358fd76153487f39d. 2021-06-02 17:13:03.173+0300 info [Task] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (5/8)#0 (d8fa853f251afd4f87776ad0b3810635) switched from CREATED to DEPLOYING. 2021-06-02 17:13:03.173+0300 info [Task] Loading JAR files for task /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (5/8)#0 (d8fa853f251afd4f87776ad0b3810635) [DEPLOYING]. 2021-06-02 17:13:03.176+0300 info [TaskSlotTableImpl] Activate slot 237b76cbdd6a12a219a9600ab9fb667a. 2021-06-02 17:13:03.176+0300 info [Task] Registering task at network: /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (5/8)#0 (d8fa853f251afd4f87776ad0b3810635) [DEPLOYING]. 2021-06-02 17:13:03.176+0300 warn [MetricGroup] The operator name Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example exceeded the 80 characters length limit and was truncated. 2021-06-02 17:13:03.177+0300 info [StreamTask] Using application-defined state backend: org.apache.flink.streaming.api.operators.sorted.state.BatchExecutionStateBackend@4b727a4d 2021-06-02 17:13:03.177+0300 info [Task] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (5/8)#0 (d8fa853f251afd4f87776ad0b3810635) switched from DEPLOYING to RUNNING. 2021-06-02 17:13:03.177+0300 warn [MetricGroup] The operator name Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example exceeded the 80 characters length limit and was truncated. 2021-06-02 17:13:03.177+0300 warn [MetricGroup] The operator name Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example exceeded the 80 characters length limit and was truncated. 2021-06-02 17:13:03.177+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (5/8) (d8fa853f251afd4f87776ad0b3810635) switched from DEPLOYING to RUNNING. 2021-06-02 17:13:03.178+0300 info [TaskExecutor] Received task /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (6/8)#0 (7e273ba7189fa887157b7cf00565aa27), deploy into slot with allocation id 237b76cbdd6a12a219a9600ab9fb667a. 2021-06-02 17:13:03.180+0300 info [Task] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (6/8)#0 (7e273ba7189fa887157b7cf00565aa27) switched from CREATED to DEPLOYING. 2021-06-02 17:13:03.180+0300 info [Task] Loading JAR files for task /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (6/8)#0 (7e273ba7189fa887157b7cf00565aa27) [DEPLOYING]. 2021-06-02 17:13:03.180+0300 info [TaskSlotTableImpl] Activate slot 5a88e10c6d023fed666e644c9af0a210. 2021-06-02 17:13:03.181+0300 info [TaskSlotTableImpl] Activate slot 96738ae5a01c330b7899819edb6918b6. 2021-06-02 17:13:03.181+0300 info [TaskSlotTableImpl] Activate slot 43d084de5b721547c53d65b821d01f6f. 2021-06-02 17:13:03.181+0300 info [TaskSlotTableImpl] Activate slot 0db08c17830fa9b358fd76153487f39d. 2021-06-02 17:13:03.182+0300 info [TaskSlotTableImpl] Activate slot 237b76cbdd6a12a219a9600ab9fb667a. 2021-06-02 17:13:03.182+0300 info [TaskSlotTableImpl] Activate slot 43d084de5b721547c53d65b821d01f6f. 2021-06-02 17:13:03.182+0300 info [TaskSlotTableImpl] Activate slot 5a88e10c6d023fed666e644c9af0a210. 2021-06-02 17:13:03.182+0300 info [TaskSlotTableImpl] Activate slot 96738ae5a01c330b7899819edb6918b6. 2021-06-02 17:13:03.182+0300 info [TaskSlotTableImpl] Activate slot 0db08c17830fa9b358fd76153487f39d. 2021-06-02 17:13:03.183+0300 info [TaskSlotTableImpl] Activate slot 237b76cbdd6a12a219a9600ab9fb667a. 2021-06-02 17:13:03.183+0300 info [Task] Registering task at network: /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (6/8)#0 (7e273ba7189fa887157b7cf00565aa27) [DEPLOYING]. 2021-06-02 17:13:03.183+0300 info [TaskSlotTableImpl] Activate slot d12f6d0d0893e1b326e2b5e94818ea06. 2021-06-02 17:13:03.183+0300 info [TaskSlotTableImpl] Activate slot 43d084de5b721547c53d65b821d01f6f. 2021-06-02 17:13:03.183+0300 info [TaskSlotTableImpl] Activate slot d12f6d0d0893e1b326e2b5e94818ea06. 2021-06-02 17:13:03.184+0300 info [TaskSlotTableImpl] Activate slot 8ab384bad8ddea1d6796c14ec8ffca36. 2021-06-02 17:13:03.184+0300 info [TaskSlotTableImpl] Activate slot 5a88e10c6d023fed666e644c9af0a210. 2021-06-02 17:13:03.184+0300 info [TaskSlotTableImpl] Activate slot 96738ae5a01c330b7899819edb6918b6. 2021-06-02 17:13:03.184+0300 info [TaskSlotTableImpl] Activate slot 0db08c17830fa9b358fd76153487f39d. 2021-06-02 17:13:03.184+0300 info [TaskSlotTableImpl] Activate slot 237b76cbdd6a12a219a9600ab9fb667a. 2021-06-02 17:13:03.184+0300 info [TaskSlotTableImpl] Activate slot d12f6d0d0893e1b326e2b5e94818ea06. 2021-06-02 17:13:03.185+0300 info [StreamTask] Using application-defined state backend: org.apache.flink.streaming.api.operators.sorted.state.BatchExecutionStateBackend@20daa6cc 2021-06-02 17:13:03.185+0300 info [Task] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (6/8)#0 (7e273ba7189fa887157b7cf00565aa27) switched from DEPLOYING to RUNNING. 2021-06-02 17:13:03.185+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (6/8) (7e273ba7189fa887157b7cf00565aa27) switched from DEPLOYING to RUNNING. 2021-06-02 17:13:03.185+0300 info [TaskExecutor] Received task /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (7/8)#0 (9b4b5cf10ffef54e899ab3c200403a3c), deploy into slot with allocation id d12f6d0d0893e1b326e2b5e94818ea06. 2021-06-02 17:13:03.185+0300 warn [MetricGroup] The operator name Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example exceeded the 80 characters length limit and was truncated. 2021-06-02 17:13:03.186+0300 info [Task] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (7/8)#0 (9b4b5cf10ffef54e899ab3c200403a3c) switched from CREATED to DEPLOYING. 2021-06-02 17:13:03.186+0300 info [Task] Loading JAR files for task /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (7/8)#0 (9b4b5cf10ffef54e899ab3c200403a3c) [DEPLOYING]. 2021-06-02 17:13:03.186+0300 info [TaskSlotTableImpl] Activate slot 8ab384bad8ddea1d6796c14ec8ffca36. 2021-06-02 17:13:03.186+0300 warn [MetricGroup] The operator name Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example exceeded the 80 characters length limit and was truncated. 2021-06-02 17:13:03.186+0300 info [Task] Registering task at network: /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (7/8)#0 (9b4b5cf10ffef54e899ab3c200403a3c) [DEPLOYING]. 2021-06-02 17:13:03.187+0300 info [StreamTask] Using application-defined state backend: org.apache.flink.streaming.api.operators.sorted.state.BatchExecutionStateBackend@d2dc15c 2021-06-02 17:13:03.187+0300 info [Task] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (7/8)#0 (9b4b5cf10ffef54e899ab3c200403a3c) switched from DEPLOYING to RUNNING. 2021-06-02 17:13:03.187+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (7/8) (9b4b5cf10ffef54e899ab3c200403a3c) switched from DEPLOYING to RUNNING. 2021-06-02 17:13:03.188+0300 info [TaskExecutor] Received task /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (8/8)#0 (ba0cbe884bf9462ac3ef39f7455be895), deploy into slot with allocation id 8ab384bad8ddea1d6796c14ec8ffca36. 2021-06-02 17:13:03.188+0300 info [Task] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (8/8)#0 (ba0cbe884bf9462ac3ef39f7455be895) switched from CREATED to DEPLOYING. 2021-06-02 17:13:03.188+0300 info [Task] Loading JAR files for task /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (8/8)#0 (ba0cbe884bf9462ac3ef39f7455be895) [DEPLOYING]. 2021-06-02 17:13:03.189+0300 warn [MetricGroup] The operator name Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example exceeded the 80 characters length limit and was truncated. 2021-06-02 17:13:03.189+0300 info [Task] Registering task at network: /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (8/8)#0 (ba0cbe884bf9462ac3ef39f7455be895) [DEPLOYING]. 2021-06-02 17:13:03.190+0300 info [StreamTask] Using application-defined state backend: org.apache.flink.streaming.api.operators.sorted.state.BatchExecutionStateBackend@221a9bc1 2021-06-02 17:13:03.190+0300 info [Task] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (8/8)#0 (ba0cbe884bf9462ac3ef39f7455be895) switched from DEPLOYING to RUNNING. 2021-06-02 17:13:03.191+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (8/8) (ba0cbe884bf9462ac3ef39f7455be895) switched from DEPLOYING to RUNNING. 2021-06-02 17:13:03.195+0300 warn [MetricGroup] The operator name Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example exceeded the 80 characters length limit and was truncated. 2021-06-02 17:13:03.198+0300 warn [MetricGroup] The operator name Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example exceeded the 80 characters length limit and was truncated. 2021-06-02 17:13:03.224+0300 info [ContinuousFileReaderOperator] No state to restore for the ContinuousFileReaderOperator (taskIdx=0). 2021-06-02 17:13:03.224+0300 info [ContinuousFileReaderOperator] No state to restore for the ContinuousFileReaderOperator (taskIdx=1). 2021-06-02 17:13:03.224+0300 info [ContinuousFileReaderOperator] No state to restore for the ContinuousFileReaderOperator (taskIdx=7). 2021-06-02 17:13:03.224+0300 info [ContinuousFileReaderOperator] No state to restore for the ContinuousFileReaderOperator (taskIdx=5). 2021-06-02 17:13:03.225+0300 info [ContinuousFileReaderOperator] No state to restore for the ContinuousFileReaderOperator (taskIdx=4). 2021-06-02 17:13:03.225+0300 info [ContinuousFileReaderOperator] No state to restore for the ContinuousFileReaderOperator (taskIdx=2). 2021-06-02 17:13:03.225+0300 info [ContinuousFileReaderOperator] No state to restore for the ContinuousFileReaderOperator (taskIdx=6). 2021-06-02 17:13:03.226+0300 info [ContinuousFileReaderOperator] No state to restore for the ContinuousFileReaderOperator (taskIdx=3). 2021-06-02 17:13:03.227+0300 info [SingleInputGate] Converting recovered input channels (1 channels) 2021-06-02 17:13:03.228+0300 info [SingleInputGate] Converting recovered input channels (1 channels) 2021-06-02 17:13:03.228+0300 info [SingleInputGate] Converting recovered input channels (1 channels) 2021-06-02 17:13:03.228+0300 info [SingleInputGate] Converting recovered input channels (1 channels) 2021-06-02 17:13:03.229+0300 info [SingleInputGate] Converting recovered input channels (1 channels) 2021-06-02 17:13:03.229+0300 info [SingleInputGate] Converting recovered input channels (1 channels) 2021-06-02 17:13:03.229+0300 info [SingleInputGate] Converting recovered input channels (1 channels) 2021-06-02 17:13:03.229+0300 info [SingleInputGate] Converting recovered input channels (1 channels) 2021-06-02 17:13:03.245+0300 info [Task] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (8/8)#0 (ba0cbe884bf9462ac3ef39f7455be895) switched from RUNNING to FINISHED. 2021-06-02 17:13:03.245+0300 info [Task] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (6/8)#0 (7e273ba7189fa887157b7cf00565aa27) switched from RUNNING to FINISHED. 2021-06-02 17:13:03.245+0300 info [Task] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (3/8)#0 (2906d5a3cbb74007054260ac3edf6385) switched from RUNNING to FINISHED. 2021-06-02 17:13:03.246+0300 info [Task] Freeing task resources for /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (6/8)#0 (7e273ba7189fa887157b7cf00565aa27). 2021-06-02 17:13:03.247+0300 info [Task] Freeing task resources for /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (3/8)#0 (2906d5a3cbb74007054260ac3edf6385). 2021-06-02 17:13:03.247+0300 info [Task] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (2/8)#0 (143e6297eb2218f5590cf068e386ba71) switched from RUNNING to FINISHED. 2021-06-02 17:13:03.247+0300 info [Task] Freeing task resources for /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (8/8)#0 (ba0cbe884bf9462ac3ef39f7455be895). 2021-06-02 17:13:03.247+0300 info [Task] Freeing task resources for /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (2/8)#0 (143e6297eb2218f5590cf068e386ba71). 2021-06-02 17:13:03.247+0300 info [Task] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (5/8)#0 (d8fa853f251afd4f87776ad0b3810635) switched from RUNNING to FINISHED. 2021-06-02 17:13:03.247+0300 info [Task] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (1/8)#0 (198c68d2f177dda7d3f1803a5c8283ca) switched from RUNNING to FINISHED. 2021-06-02 17:13:03.248+0300 info [Task] Freeing task resources for /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (5/8)#0 (d8fa853f251afd4f87776ad0b3810635). 2021-06-02 17:13:03.248+0300 info [Task] Freeing task resources for /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (1/8)#0 (198c68d2f177dda7d3f1803a5c8283ca). 2021-06-02 17:13:03.248+0300 info [TaskExecutor] Un-registering task and sending final execution state FINISHED to JobManager for task /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (6/8)#0 7e273ba7189fa887157b7cf00565aa27. 2021-06-02 17:13:03.248+0300 info [TaskExecutor] Un-registering task and sending final execution state FINISHED to JobManager for task /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (2/8)#0 143e6297eb2218f5590cf068e386ba71. 2021-06-02 17:13:03.248+0300 info [TaskExecutor] Un-registering task and sending final execution state FINISHED to JobManager for task /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (3/8)#0 2906d5a3cbb74007054260ac3edf6385. 2021-06-02 17:13:03.249+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (6/8) (7e273ba7189fa887157b7cf00565aa27) switched from RUNNING to FINISHED. 2021-06-02 17:13:03.249+0300 info [TaskExecutor] Un-registering task and sending final execution state FINISHED to JobManager for task /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (8/8)#0 ba0cbe884bf9462ac3ef39f7455be895. 2021-06-02 17:13:03.249+0300 info [TaskExecutor] Un-registering task and sending final execution state FINISHED to JobManager for task /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (1/8)#0 198c68d2f177dda7d3f1803a5c8283ca. 2021-06-02 17:13:03.249+0300 info [TaskExecutor] Un-registering task and sending final execution state FINISHED to JobManager for task /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (5/8)#0 d8fa853f251afd4f87776ad0b3810635. 2021-06-02 17:13:03.249+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (2/8) (143e6297eb2218f5590cf068e386ba71) switched from RUNNING to FINISHED. 2021-06-02 17:13:03.250+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (3/8) (2906d5a3cbb74007054260ac3edf6385) switched from RUNNING to FINISHED. 2021-06-02 17:13:03.251+0300 info [Task] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (4/8)#0 (1c884a8478c56d4a9270d30c140724e8) switched from RUNNING to FINISHED. 2021-06-02 17:13:03.251+0300 info [Task] Freeing task resources for /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (4/8)#0 (1c884a8478c56d4a9270d30c140724e8). 2021-06-02 17:13:03.251+0300 info [TaskExecutor] Un-registering task and sending final execution state FINISHED to JobManager for task /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (4/8)#0 1c884a8478c56d4a9270d30c140724e8. 2021-06-02 17:13:03.252+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (8/8) (ba0cbe884bf9462ac3ef39f7455be895) switched from RUNNING to FINISHED. 2021-06-02 17:13:03.253+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (1/8) (198c68d2f177dda7d3f1803a5c8283ca) switched from RUNNING to FINISHED. 2021-06-02 17:13:03.254+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (5/8) (d8fa853f251afd4f87776ad0b3810635) switched from RUNNING to FINISHED. 2021-06-02 17:13:03.256+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (4/8) (1c884a8478c56d4a9270d30c140724e8) switched from RUNNING to FINISHED. 2021-06-02 17:13:03.551+0300 warn [NativeCodeLoader] Unable to load native-hadoop library for your platform... using builtin-java classes where applicable 2021-06-02 17:13:04.266+0300 error [ParquetPojoInputFormat] Fields number is 6. 2021-06-02 17:13:04.329+0300 info [CodecPool] Got brand-new decompressor [.gz] 2021-06-02 17:13:04.618+0300 info [CodecPool] Got brand-new compressor [.gz] 2021-06-02 17:13:13.456+0300 warn [ContinuousFileReaderOperator] not processing any records while closed 2021-06-02 17:13:13.458+0300 warn [Task] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (7/8)#0 (9b4b5cf10ffef54e899ab3c200403a3c) switched from RUNNING to FAILED. java.nio.channels.ClosedChannelException at java.base/sun.nio.ch.FileChannelImpl.ensureOpen(FileChannelImpl.java:156) at java.base/sun.nio.ch.FileChannelImpl.position(FileChannelImpl.java:331) at org.apache.flink.core.fs.local.LocalRecoverableFsDataOutputStream.getPos(LocalRecoverableFsDataOutputStream.java:103) at org.apache.flink.formats.parquet.PositionOutputStreamAdapter.getPos(PositionOutputStreamAdapter.java:48) at org.apache.parquet.hadoop.ParquetFileWriter.startBlock(ParquetFileWriter.java:349) at org.apache.parquet.hadoop.InternalParquetRecordWriter.flushRowGroupToStore(InternalParquetRecordWriter.java:171) at org.apache.parquet.hadoop.InternalParquetRecordWriter.close(InternalParquetRecordWriter.java:114) at org.apache.parquet.hadoop.ParquetWriter.close(ParquetWriter.java:310) at org.apache.flink.formats.parquet.ParquetBulkWriter.finish(ParquetBulkWriter.java:62) at org.apache.flink.streaming.api.functions.sink.filesystem.BulkPartWriter.closeForCommit(BulkPartWriter.java:61) at org.apache.flink.connector.file.sink.writer.FileWriterBucket.closePartFile(FileWriterBucket.java:279) at org.apache.flink.connector.file.sink.writer.FileWriterBucket.prepareCommit(FileWriterBucket.java:201) at org.apache.flink.connector.file.sink.writer.FileWriter.prepareCommit(FileWriter.java:200) at org.apache.flink.streaming.runtime.operators.sink.AbstractSinkWriterOperator.endInput(AbstractSinkWriterOperator.java:97) at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.endOperatorInput(StreamOperatorWrapper.java:91) at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.lambda$close$0(StreamOperatorWrapper.java:128) at org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$1.runThrowing(StreamTaskActionExecutor.java:50) at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.close(StreamOperatorWrapper.java:128) at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.close(StreamOperatorWrapper.java:135) at org.apache.flink.streaming.runtime.tasks.OperatorChain.closeOperators(OperatorChain.java:439) at org.apache.flink.streaming.runtime.tasks.StreamTask.afterInvoke(StreamTask.java:629) at org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:591) at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:758) at org.apache.flink.runtime.taskmanager.Task.run(Task.java:573) 2021-06-02 17:13:13.465+0300 info [Task] Freeing task resources for /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (7/8)#0 (9b4b5cf10ffef54e899ab3c200403a3c). 2021-06-02 17:13:13.466+0300 info [TaskExecutor] Un-registering task and sending final execution state FAILED to JobManager for task /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (7/8)#0 9b4b5cf10ffef54e899ab3c200403a3c. 2021-06-02 17:13:13.466+0300 info [ExecutionGraph] /Users/user/compacter/tmp//parquet_example -> Sink Writer: /Users/user/compacter/tmp/tttt//parquet_example (7/8) (9b4b5cf10ffef54e899ab3c200403a3c) switched from RUNNING to FAILED on 264008fd-b7b7-4639-a3d4-ae430793106a @ localhost (dataPort=-1). java.nio.channels.ClosedChannelException at java.base/sun.nio.ch.FileChannelImpl.ensureOpen(FileChannelImpl.java:156) at java.base/sun.nio.ch.FileChannelImpl.position(FileChannelImpl.java:331) at org.apache.flink.core.fs.local.LocalRecoverableFsDataOutputStream.getPos(LocalRecoverableFsDataOutputStream.java:103) at org.apache.flink.formats.parquet.PositionOutputStreamAdapter.getPos(PositionOutputStreamAdapter.java:48) at org.apache.parquet.hadoop.ParquetFileWriter.startBlock(ParquetFileWriter.java:349) at org.apache.parquet.hadoop.InternalParquetRecordWriter.flushRowGroupToStore(InternalParquetRecordWriter.java:171) at org.apache.parquet.hadoop.InternalParquetRecordWriter.close(InternalParquetRecordWriter.java:114) at org.apache.parquet.hadoop.ParquetWriter.close(ParquetWriter.java:310) at org.apache.flink.formats.parquet.ParquetBulkWriter.finish(ParquetBulkWriter.java:62) at org.apache.flink.streaming.api.functions.sink.filesystem.BulkPartWriter.closeForCommit(BulkPartWriter.java:61) at org.apache.flink.connector.file.sink.writer.FileWriterBucket.closePartFile(FileWriterBucket.java:279) at org.apache.flink.connector.file.sink.writer.FileWriterBucket.prepareCommit(FileWriterBucket.java:201) at org.apache.flink.connector.file.sink.writer.FileWriter.prepareCommit(FileWriter.java:200) at org.apache.flink.streaming.runtime.operators.sink.AbstractSinkWriterOperator.endInput(AbstractSinkWriterOperator.java:97) at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.endOperatorInput(StreamOperatorWrapper.java:91) at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.lambda$close$0(StreamOperatorWrapper.java:128) at org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$1.runThrowing(StreamTaskActionExecutor.java:50) at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.close(StreamOperatorWrapper.java:128) at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.close(StreamOperatorWrapper.java:135) at org.apache.flink.streaming.runtime.tasks.OperatorChain.closeOperators(OperatorChain.java:439) at org.apache.flink.streaming.runtime.tasks.StreamTask.afterInvoke(StreamTask.java:629) at org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:591) at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:758) at org.apache.flink.runtime.taskmanager.Task.run(Task.java:573) 2021-06-02 17:13:13.468+0300 info [RestartPipelinedRegionFailoverStrategy] Calculating tasks to restart to recover the failed task 20ba6b65f97481d5570070de90e4e791_6. 2021-06-02 17:13:13.469+0300 info [RestartPipelinedRegionFailoverStrategy] 1 tasks should be restarted to recover the failed task 20ba6b65f97481d5570070de90e4e791_6. 2021-06-02 17:13:13.470+0300 info [ExecutionGraph] Job hdfs_parquet_compacter (d70f7e1ff824fb485de8a4dbf2b485b0) switched from state RUNNING to FAILING. org.apache.flink.runtime.JobException: Recovery is suppressed by NoRestartBackoffTimeStrategy at org.apache.flink.runtime.executiongraph.failover.flip1.ExecutionFailureHandler.handleFailure(ExecutionFailureHandler.java:118) at org.apache.flink.runtime.executiongraph.failover.flip1.ExecutionFailureHandler.getFailureHandlingResult(ExecutionFailureHandler.java:80) at org.apache.flink.runtime.scheduler.DefaultScheduler.handleTaskFailure(DefaultScheduler.java:233) at org.apache.flink.runtime.scheduler.DefaultScheduler.maybeHandleTaskFailure(DefaultScheduler.java:224) at org.apache.flink.runtime.scheduler.DefaultScheduler.updateTaskExecutionStateInternal(DefaultScheduler.java:215) at org.apache.flink.runtime.scheduler.SchedulerBase.updateTaskExecutionState(SchedulerBase.java:666) at org.apache.flink.runtime.scheduler.SchedulerNG.updateTaskExecutionState(SchedulerNG.java:89) at org.apache.flink.runtime.jobmaster.JobMaster.updateTaskExecutionState(JobMaster.java:446) at jdk.internal.reflect.GeneratedMethodAccessor18.invoke(Unknown Source) at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.base/java.lang.reflect.Method.invoke(Method.java:564) at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcInvocation(AkkaRpcActor.java:305) at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcMessage(AkkaRpcActor.java:212) at org.apache.flink.runtime.rpc.akka.FencedAkkaRpcActor.handleRpcMessage(FencedAkkaRpcActor.java:77) at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleMessage(AkkaRpcActor.java:158) at akka.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:26) at akka.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:21) at scala.PartialFunction.applyOrElse(PartialFunction.scala:127) at scala.PartialFunction.applyOrElse$(PartialFunction.scala:126) at akka.japi.pf.UnitCaseStatement.applyOrElse(CaseStatements.scala:21) at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:175) at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:176) at akka.actor.Actor.aroundReceive(Actor.scala:517) at akka.actor.Actor.aroundReceive$(Actor.scala:515) at akka.actor.AbstractActor.aroundReceive(AbstractActor.scala:225) at akka.actor.ActorCell.receiveMessage(ActorCell.scala:592) at akka.actor.ActorCell.invoke(ActorCell.scala:561) at akka.dispatch.Mailbox.processMailbox(Mailbox.scala:258) at akka.dispatch.Mailbox.run(Mailbox.scala:225) at akka.dispatch.Mailbox.exec(Mailbox.scala:235) at akka.dispatch.forkjoin.ForkJoinTask.doExec(ForkJoinTask.java:260) at akka.dispatch.forkjoin.ForkJoinPool$WorkQueue.runTask(ForkJoinPool.java:1339) at akka.dispatch.forkjoin.ForkJoinPool.runWorker(ForkJoinPool.java:1979) at akka.dispatch.forkjoin.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:107) Caused by: java.nio.channels.ClosedChannelException at java.base/sun.nio.ch.FileChannelImpl.ensureOpen(FileChannelImpl.java:156) at java.base/sun.nio.ch.FileChannelImpl.position(FileChannelImpl.java:331) at org.apache.flink.core.fs.local.LocalRecoverableFsDataOutputStream.getPos(LocalRecoverableFsDataOutputStream.java:103) at org.apache.flink.formats.parquet.PositionOutputStreamAdapter.getPos(PositionOutputStreamAdapter.java:48) at org.apache.parquet.hadoop.ParquetFileWriter.startBlock(ParquetFileWriter.java:349) at org.apache.parquet.hadoop.InternalParquetRecordWriter.flushRowGroupToStore(InternalParquetRecordWriter.java:171) at org.apache.parquet.hadoop.InternalParquetRecordWriter.close(InternalParquetRecordWriter.java:114) at org.apache.parquet.hadoop.ParquetWriter.close(ParquetWriter.java:310) at org.apache.flink.formats.parquet.ParquetBulkWriter.finish(ParquetBulkWriter.java:62) at org.apache.flink.streaming.api.functions.sink.filesystem.BulkPartWriter.closeForCommit(BulkPartWriter.java:61) at org.apache.flink.connector.file.sink.writer.FileWriterBucket.closePartFile(FileWriterBucket.java:279) at org.apache.flink.connector.file.sink.writer.FileWriterBucket.prepareCommit(FileWriterBucket.java:201) at org.apache.flink.connector.file.sink.writer.FileWriter.prepareCommit(FileWriter.java:200) at org.apache.flink.streaming.runtime.operators.sink.AbstractSinkWriterOperator.endInput(AbstractSinkWriterOperator.java:97) at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.endOperatorInput(StreamOperatorWrapper.java:91) at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.lambda$close$0(StreamOperatorWrapper.java:128) at org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$1.runThrowing(StreamTaskActionExecutor.java:50) at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.close(StreamOperatorWrapper.java:128) at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.close(StreamOperatorWrapper.java:135) at org.apache.flink.streaming.runtime.tasks.OperatorChain.closeOperators(OperatorChain.java:439) at org.apache.flink.streaming.runtime.tasks.StreamTask.afterInvoke(StreamTask.java:629) at org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:591) at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:758) at org.apache.flink.runtime.taskmanager.Task.run(Task.java:573) 2021-06-02 17:13:13.477+0300 info [ExecutionGraph] Discarding the results produced by task execution 6330e22e7399f04ba7087e1b2e885044. 2021-06-02 17:13:13.477+0300 info [ExecutionGraph] Discarding the results produced by task execution 198c68d2f177dda7d3f1803a5c8283ca. 2021-06-02 17:13:13.477+0300 info [ExecutionGraph] Discarding the results produced by task execution 143e6297eb2218f5590cf068e386ba71. 2021-06-02 17:13:13.477+0300 info [ExecutionGraph] Discarding the results produced by task execution 2906d5a3cbb74007054260ac3edf6385. 2021-06-02 17:13:13.477+0300 info [ExecutionGraph] Discarding the results produced by task execution 1c884a8478c56d4a9270d30c140724e8. 2021-06-02 17:13:13.477+0300 info [ExecutionGraph] Discarding the results produced by task execution d8fa853f251afd4f87776ad0b3810635. 2021-06-02 17:13:13.477+0300 info [ExecutionGraph] Discarding the results produced by task execution 7e273ba7189fa887157b7cf00565aa27. 2021-06-02 17:13:13.477+0300 info [ExecutionGraph] Discarding the results produced by task execution ba0cbe884bf9462ac3ef39f7455be895. 2021-06-02 17:13:13.477+0300 info [ExecutionGraph] Sink Committer: /Users/user/compacter/tmp/tttt//parquet_example (1/1) (045d421cb9c1085378b720c8944d2c20) switched from CREATED to CANCELING. 2021-06-02 17:13:13.478+0300 info [ExecutionGraph] Sink Committer: /Users/user/compacter/tmp/tttt//parquet_example (1/1) (045d421cb9c1085378b720c8944d2c20) switched from CANCELING to CANCELED. 2021-06-02 17:13:13.478+0300 info [ExecutionGraph] Discarding the results produced by task execution 045d421cb9c1085378b720c8944d2c20. 2021-06-02 17:13:13.478+0300 info [ExecutionGraph] Job hdfs_parquet_compacter (d70f7e1ff824fb485de8a4dbf2b485b0) switched from state FAILING to FAILED. org.apache.flink.runtime.JobException: Recovery is suppressed by NoRestartBackoffTimeStrategy at org.apache.flink.runtime.executiongraph.failover.flip1.ExecutionFailureHandler.handleFailure(ExecutionFailureHandler.java:118) at org.apache.flink.runtime.executiongraph.failover.flip1.ExecutionFailureHandler.getFailureHandlingResult(ExecutionFailureHandler.java:80) at org.apache.flink.runtime.scheduler.DefaultScheduler.handleTaskFailure(DefaultScheduler.java:233) at org.apache.flink.runtime.scheduler.DefaultScheduler.maybeHandleTaskFailure(DefaultScheduler.java:224) at org.apache.flink.runtime.scheduler.DefaultScheduler.updateTaskExecutionStateInternal(DefaultScheduler.java:215) at org.apache.flink.runtime.scheduler.SchedulerBase.updateTaskExecutionState(SchedulerBase.java:666) at org.apache.flink.runtime.scheduler.SchedulerNG.updateTaskExecutionState(SchedulerNG.java:89) at org.apache.flink.runtime.jobmaster.JobMaster.updateTaskExecutionState(JobMaster.java:446) at jdk.internal.reflect.GeneratedMethodAccessor18.invoke(Unknown Source) at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.base/java.lang.reflect.Method.invoke(Method.java:564) at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcInvocation(AkkaRpcActor.java:305) at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcMessage(AkkaRpcActor.java:212) at org.apache.flink.runtime.rpc.akka.FencedAkkaRpcActor.handleRpcMessage(FencedAkkaRpcActor.java:77) at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleMessage(AkkaRpcActor.java:158) at akka.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:26) at akka.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:21) at scala.PartialFunction.applyOrElse(PartialFunction.scala:127) at scala.PartialFunction.applyOrElse$(PartialFunction.scala:126) at akka.japi.pf.UnitCaseStatement.applyOrElse(CaseStatements.scala:21) at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:175) at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:176) at akka.actor.Actor.aroundReceive(Actor.scala:517) at akka.actor.Actor.aroundReceive$(Actor.scala:515) at akka.actor.AbstractActor.aroundReceive(AbstractActor.scala:225) at akka.actor.ActorCell.receiveMessage(ActorCell.scala:592) at akka.actor.ActorCell.invoke(ActorCell.scala:561) at akka.dispatch.Mailbox.processMailbox(Mailbox.scala:258) at akka.dispatch.Mailbox.run(Mailbox.scala:225) at akka.dispatch.Mailbox.exec(Mailbox.scala:235) at akka.dispatch.forkjoin.ForkJoinTask.doExec(ForkJoinTask.java:260) at akka.dispatch.forkjoin.ForkJoinPool$WorkQueue.runTask(ForkJoinPool.java:1339) at akka.dispatch.forkjoin.ForkJoinPool.runWorker(ForkJoinPool.java:1979) at akka.dispatch.forkjoin.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:107) Caused by: java.nio.channels.ClosedChannelException at java.base/sun.nio.ch.FileChannelImpl.ensureOpen(FileChannelImpl.java:156) at java.base/sun.nio.ch.FileChannelImpl.position(FileChannelImpl.java:331) at org.apache.flink.core.fs.local.LocalRecoverableFsDataOutputStream.getPos(LocalRecoverableFsDataOutputStream.java:103) at org.apache.flink.formats.parquet.PositionOutputStreamAdapter.getPos(PositionOutputStreamAdapter.java:48) at org.apache.parquet.hadoop.ParquetFileWriter.startBlock(ParquetFileWriter.java:349) at org.apache.parquet.hadoop.InternalParquetRecordWriter.flushRowGroupToStore(InternalParquetRecordWriter.java:171) at org.apache.parquet.hadoop.InternalParquetRecordWriter.close(InternalParquetRecordWriter.java:114) at org.apache.parquet.hadoop.ParquetWriter.close(ParquetWriter.java:310) at org.apache.flink.formats.parquet.ParquetBulkWriter.finish(ParquetBulkWriter.java:62) at org.apache.flink.streaming.api.functions.sink.filesystem.BulkPartWriter.closeForCommit(BulkPartWriter.java:61) at org.apache.flink.connector.file.sink.writer.FileWriterBucket.closePartFile(FileWriterBucket.java:279) at org.apache.flink.connector.file.sink.writer.FileWriterBucket.prepareCommit(FileWriterBucket.java:201) at org.apache.flink.connector.file.sink.writer.FileWriter.prepareCommit(FileWriter.java:200) at org.apache.flink.streaming.runtime.operators.sink.AbstractSinkWriterOperator.endInput(AbstractSinkWriterOperator.java:97) at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.endOperatorInput(StreamOperatorWrapper.java:91) at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.lambda$close$0(StreamOperatorWrapper.java:128) at org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$1.runThrowing(StreamTaskActionExecutor.java:50) at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.close(StreamOperatorWrapper.java:128) at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.close(StreamOperatorWrapper.java:135) at org.apache.flink.streaming.runtime.tasks.OperatorChain.closeOperators(OperatorChain.java:439) at org.apache.flink.streaming.runtime.tasks.StreamTask.afterInvoke(StreamTask.java:629) at org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:591) at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:758) at org.apache.flink.runtime.taskmanager.Task.run(Task.java:573) 2021-06-02 17:13:13.482+0300 info [CheckpointCoordinator] Stopping checkpoint coordinator for job d70f7e1ff824fb485de8a4dbf2b485b0. 2021-06-02 17:13:13.482+0300 info [StandaloneCompletedCheckpointStore] Shutting down 2021-06-02 17:13:13.486+0300 info [StandaloneDispatcher] Job d70f7e1ff824fb485de8a4dbf2b485b0 reached globally terminal state FAILED. 2021-06-02 17:13:13.486+0300 info [MiniCluster] Shutting down Flink Mini Cluster 2021-06-02 17:13:13.486+0300 info [TaskExecutor] Stopping TaskExecutor akka://flink/user/rpc/taskmanager_0. 2021-06-02 17:13:13.486+0300 info [TaskExecutor] Close ResourceManager connection a40c5bbcceefd20258e8ba206570d75e. 2021-06-02 17:13:13.487+0300 info [DispatcherRestEndpoint] Shutting down rest endpoint. 2021-06-02 17:13:13.487+0300 info [StandaloneResourceManager] Closing TaskExecutor connection 264008fd-b7b7-4639-a3d4-ae430793106a because: The TaskExecutor is shutting down. 2021-06-02 17:13:13.488+0300 info [TaskExecutor] Close JobManager connection for job d70f7e1ff824fb485de8a4dbf2b485b0. 2021-06-02 17:13:13.489+0300 info [JobMaster] Stopping the JobMaster for job hdfs_parquet_compacter(d70f7e1ff824fb485de8a4dbf2b485b0). 2021-06-02 17:13:13.489+0300 info [SlotPoolImpl] Suspending SlotPool. 2021-06-02 17:13:13.489+0300 info [JobMaster] Close ResourceManager connection a40c5bbcceefd20258e8ba206570d75e: Stopping JobMaster for job hdfs_parquet_compacter(d70f7e1ff824fb485de8a4dbf2b485b0).. 2021-06-02 17:13:13.490+0300 info [SlotPoolImpl] Stopping SlotPool. 2021-06-02 17:13:13.490+0300 info [TaskSlotTableImpl] Free slot TaskSlot(index:0, state:ALLOCATED, resource profile: ResourceProfile{managedMemory=16.000mb (16777216 bytes), networkMemory=8.000mb (8388608 bytes)}, allocationId: 2b3f94c62b0602f21dd1b4af99bec0fc, jobId: d70f7e1ff824fb485de8a4dbf2b485b0). 2021-06-02 17:13:13.490+0300 info [StandaloneResourceManager] Disconnect job manager bc024d08781064bc55764aaea3b9435d@akka://flink/user/rpc/jobmanager_3 for job d70f7e1ff824fb485de8a4dbf2b485b0 from the resource manager. 2021-06-02 17:13:13.491+0300 info [TaskSlotTableImpl] Free slot TaskSlot(index:3, state:ALLOCATED, resource profile: ResourceProfile{managedMemory=16.000mb (16777216 bytes), networkMemory=8.000mb (8388608 bytes)}, allocationId: 43d084de5b721547c53d65b821d01f6f, jobId: d70f7e1ff824fb485de8a4dbf2b485b0). 2021-06-02 17:13:13.491+0300 info [TaskSlotTableImpl] Free slot TaskSlot(index:6, state:ALLOCATED, resource profile: ResourceProfile{managedMemory=16.000mb (16777216 bytes), networkMemory=8.000mb (8388608 bytes)}, allocationId: d12f6d0d0893e1b326e2b5e94818ea06, jobId: d70f7e1ff824fb485de8a4dbf2b485b0). 2021-06-02 17:13:13.492+0300 info [TaskSlotTableImpl] Free slot TaskSlot(index:7, state:ALLOCATED, resource profile: ResourceProfile{managedMemory=16.000mb (16777216 bytes), networkMemory=8.000mb (8388608 bytes)}, allocationId: 8ab384bad8ddea1d6796c14ec8ffca36, jobId: d70f7e1ff824fb485de8a4dbf2b485b0). 2021-06-02 17:13:13.492+0300 info [TaskSlotTableImpl] Free slot TaskSlot(index:1, state:ALLOCATED, resource profile: ResourceProfile{managedMemory=16.000mb (16777216 bytes), networkMemory=8.000mb (8388608 bytes)}, allocationId: 5a88e10c6d023fed666e644c9af0a210, jobId: d70f7e1ff824fb485de8a4dbf2b485b0). 2021-06-02 17:13:13.493+0300 info [TaskSlotTableImpl] Free slot TaskSlot(index:2, state:ALLOCATED, resource profile: ResourceProfile{managedMemory=16.000mb (16777216 bytes), networkMemory=8.000mb (8388608 bytes)}, allocationId: 96738ae5a01c330b7899819edb6918b6, jobId: d70f7e1ff824fb485de8a4dbf2b485b0). 2021-06-02 17:13:13.493+0300 info [TaskSlotTableImpl] Free slot TaskSlot(index:4, state:ALLOCATED, resource profile: ResourceProfile{managedMemory=16.000mb (16777216 bytes), networkMemory=8.000mb (8388608 bytes)}, allocationId: 0db08c17830fa9b358fd76153487f39d, jobId: d70f7e1ff824fb485de8a4dbf2b485b0). 2021-06-02 17:13:13.494+0300 info [TaskSlotTableImpl] Free slot TaskSlot(index:5, state:ALLOCATED, resource profile: ResourceProfile{managedMemory=16.000mb (16777216 bytes), networkMemory=8.000mb (8388608 bytes)}, allocationId: 237b76cbdd6a12a219a9600ab9fb667a, jobId: d70f7e1ff824fb485de8a4dbf2b485b0). 2021-06-02 17:13:13.501+0300 info [TaskExecutor] JobManager for job d70f7e1ff824fb485de8a4dbf2b485b0 with leader id bc024d08781064bc55764aaea3b9435d lost leadership. 2021-06-02 17:13:13.503+0300 info [DispatcherRestEndpoint] Removing cache directory /var/folders/ms/18j3nxjd6rzcj9mmgly3n0rrq_g_vb/T/flink-web-ui 2021-06-02 17:13:13.504+0300 error [Main] org.apache.flink.runtime.client.JobExecutionException: Job execution failed. at org.apache.flink.runtime.jobmaster.JobResult.toJobExecutionResult(JobResult.java:144) at org.apache.flink.runtime.minicluster.MiniClusterJobClient.lambda$getJobExecutionResult$2(MiniClusterJobClient.java:117) at java.base/java.util.concurrent.CompletableFuture$UniApply.tryFire(CompletableFuture.java:642) at java.base/java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:506) at java.base/java.util.concurrent.CompletableFuture.complete(CompletableFuture.java:2137) at org.apache.flink.runtime.rpc.akka.AkkaInvocationHandler.lambda$invokeRpc$0(AkkaInvocationHandler.java:237) at java.base/java.util.concurrent.CompletableFuture.uniWhenComplete(CompletableFuture.java:859) at java.base/java.util.concurrent.CompletableFuture$UniWhenComplete.tryFire(CompletableFuture.java:837) at java.base/java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:506) at java.base/java.util.concurrent.CompletableFuture.complete(CompletableFuture.java:2137) at org.apache.flink.runtime.concurrent.FutureUtils$1.onComplete(FutureUtils.java:1061) at akka.dispatch.OnComplete.internal(Future.scala:264) at akka.dispatch.OnComplete.internal(Future.scala:261) at akka.dispatch.japi$CallbackBridge.apply(Future.scala:191) at akka.dispatch.japi$CallbackBridge.apply(Future.scala:188) at scala.concurrent.impl.CallbackRunnable.run(Promise.scala:64) at org.apache.flink.runtime.concurrent.Executors$DirectExecutionContext.execute(Executors.java:73) at scala.concurrent.impl.CallbackRunnable.executeWithValue(Promise.scala:72) at scala.concurrent.impl.Promise$DefaultPromise.$anonfun$tryComplete$1(Promise.scala:288) at scala.concurrent.impl.Promise$DefaultPromise.$anonfun$tryComplete$1$adapted(Promise.scala:288) at scala.concurrent.impl.Promise$DefaultPromise.tryComplete(Promise.scala:288) at akka.pattern.PromiseActorRef.$bang(AskSupport.scala:573) at akka.pattern.PipeToSupport$PipeableFuture$$anonfun$pipeTo$1.applyOrElse(PipeToSupport.scala:22) at akka.pattern.PipeToSupport$PipeableFuture$$anonfun$pipeTo$1.applyOrElse(PipeToSupport.scala:21) at scala.concurrent.Future.$anonfun$andThen$1(Future.scala:536) at scala.concurrent.impl.Promise.liftedTree1$1(Promise.scala:33) at scala.concurrent.impl.Promise.$anonfun$transform$1(Promise.scala:33) at scala.concurrent.impl.CallbackRunnable.run(Promise.scala:64) at akka.dispatch.BatchingExecutor$AbstractBatch.processBatch(BatchingExecutor.scala:55) at akka.dispatch.BatchingExecutor$BlockableBatch.$anonfun$run$1(BatchingExecutor.scala:91) at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) at scala.concurrent.BlockContext$.withBlockContext(BlockContext.scala:85) at akka.dispatch.BatchingExecutor$BlockableBatch.run(BatchingExecutor.scala:91) at akka.dispatch.TaskInvocation.run(AbstractDispatcher.scala:40) at akka.dispatch.ForkJoinExecutorConfigurator$AkkaForkJoinTask.exec(ForkJoinExecutorConfigurator.scala:44) at akka.dispatch.forkjoin.ForkJoinTask.doExec(ForkJoinTask.java:260) at akka.dispatch.forkjoin.ForkJoinPool$WorkQueue.runTask(ForkJoinPool.java:1339) at akka.dispatch.forkjoin.ForkJoinPool.runWorker(ForkJoinPool.java:1979) at akka.dispatch.forkjoin.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:107) Caused by: org.apache.flink.runtime.JobException: Recovery is suppressed by NoRestartBackoffTimeStrategy at org.apache.flink.runtime.executiongraph.failover.flip1.ExecutionFailureHandler.handleFailure(ExecutionFailureHandler.java:118) at org.apache.flink.runtime.executiongraph.failover.flip1.ExecutionFailureHandler.getFailureHandlingResult(ExecutionFailureHandler.java:80) at org.apache.flink.runtime.scheduler.DefaultScheduler.handleTaskFailure(DefaultScheduler.java:233) at org.apache.flink.runtime.scheduler.DefaultScheduler.maybeHandleTaskFailure(DefaultScheduler.java:224) at org.apache.flink.runtime.scheduler.DefaultScheduler.updateTaskExecutionStateInternal(DefaultScheduler.java:215) at org.apache.flink.runtime.scheduler.SchedulerBase.updateTaskExecutionState(SchedulerBase.java:666) at org.apache.flink.runtime.scheduler.SchedulerNG.updateTaskExecutionState(SchedulerNG.java:89) at org.apache.flink.runtime.jobmaster.JobMaster.updateTaskExecutionState(JobMaster.java:446) at jdk.internal.reflect.GeneratedMethodAccessor18.invoke(Unknown Source) at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.base/java.lang.reflect.Method.invoke(Method.java:564) at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcInvocation(AkkaRpcActor.java:305) at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcMessage(AkkaRpcActor.java:212) at org.apache.flink.runtime.rpc.akka.FencedAkkaRpcActor.handleRpcMessage(FencedAkkaRpcActor.java:77) at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleMessage(AkkaRpcActor.java:158) at akka.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:26) at akka.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:21) at scala.PartialFunction.applyOrElse(PartialFunction.scala:127) at scala.PartialFunction.applyOrElse$(PartialFunction.scala:126) at akka.japi.pf.UnitCaseStatement.applyOrElse(CaseStatements.scala:21) at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:175) at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:176) at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:176) at akka.actor.Actor.aroundReceive(Actor.scala:517) at akka.actor.Actor.aroundReceive$(Actor.scala:515) at akka.actor.AbstractActor.aroundReceive(AbstractActor.scala:225) at akka.actor.ActorCell.receiveMessage(ActorCell.scala:592) at akka.actor.ActorCell.invoke(ActorCell.scala:561) at akka.dispatch.Mailbox.processMailbox(Mailbox.scala:258) at akka.dispatch.Mailbox.run(Mailbox.scala:225) at akka.dispatch.Mailbox.exec(Mailbox.scala:235) ... 4 more Caused by: java.nio.channels.ClosedChannelException at java.base/sun.nio.ch.FileChannelImpl.ensureOpen(FileChannelImpl.java:156) at java.base/sun.nio.ch.FileChannelImpl.position(FileChannelImpl.java:331) at org.apache.flink.core.fs.local.LocalRecoverableFsDataOutputStream.getPos(LocalRecoverableFsDataOutputStream.java:103) at org.apache.flink.formats.parquet.PositionOutputStreamAdapter.getPos(PositionOutputStreamAdapter.java:48) at org.apache.parquet.hadoop.ParquetFileWriter.startBlock(ParquetFileWriter.java:349) at org.apache.parquet.hadoop.InternalParquetRecordWriter.flushRowGroupToStore(InternalParquetRecordWriter.java:171) at org.apache.parquet.hadoop.InternalParquetRecordWriter.close(InternalParquetRecordWriter.java:114) at org.apache.parquet.hadoop.ParquetWriter.close(ParquetWriter.java:310) at org.apache.flink.formats.parquet.ParquetBulkWriter.finish(ParquetBulkWriter.java:62) at org.apache.flink.streaming.api.functions.sink.filesystem.BulkPartWriter.closeForCommit(BulkPartWriter.java:61) at org.apache.flink.connector.file.sink.writer.FileWriterBucket.closePartFile(FileWriterBucket.java:279) at org.apache.flink.connector.file.sink.writer.FileWriterBucket.prepareCommit(FileWriterBucket.java:201) at org.apache.flink.connector.file.sink.writer.FileWriter.prepareCommit(FileWriter.java:200) at org.apache.flink.streaming.runtime.operators.sink.AbstractSinkWriterOperator.endInput(AbstractSinkWriterOperator.java:97) at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.endOperatorInput(StreamOperatorWrapper.java:91) at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.lambda$close$0(StreamOperatorWrapper.java:128) at org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$1.runThrowing(StreamTaskActionExecutor.java:50) at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.close(StreamOperatorWrapper.java:128) at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.close(StreamOperatorWrapper.java:135) at org.apache.flink.streaming.runtime.tasks.OperatorChain.closeOperators(OperatorChain.java:439) at org.apache.flink.streaming.runtime.tasks.StreamTask.afterInvoke(StreamTask.java:629) at org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:591) at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:758) at org.apache.flink.runtime.taskmanager.Task.run(Task.java:573) at java.base/java.lang.Thread.run(Thread.java:832) - (Main.scala:50) 2021-06-02 17:13:13.506+0300 info [DispatcherRestEndpoint] Shut down complete. 2021-06-02 17:13:13.507+0300 info [StandaloneResourceManager] Shut down cluster because application is in CANCELED, diagnostics DispatcherResourceManagerComponent has been closed.. 2021-06-02 17:13:13.508+0300 info [DispatcherResourceManagerComponent] Closing components. 2021-06-02 17:13:13.508+0300 info [SessionDispatcherLeaderProcess] Stopping SessionDispatcherLeaderProcess. 2021-06-02 17:13:13.509+0300 info [StandaloneDispatcher] Stopping dispatcher akka://flink/user/rpc/dispatcher_2. 2021-06-02 17:13:13.509+0300 info [StandaloneDispatcher] Stopping all currently running jobs of dispatcher akka://flink/user/rpc/dispatcher_2. 2021-06-02 17:13:13.509+0300 info [SlotManagerImpl] Closing the SlotManager. 2021-06-02 17:13:13.509+0300 info [SlotManagerImpl] Suspending the SlotManager. 2021-06-02 17:13:13.509+0300 info [BackPressureRequestCoordinator] Shutting down back pressure request coordinator. 2021-06-02 17:13:13.510+0300 info [StandaloneDispatcher] Stopped dispatcher akka://flink/user/rpc/dispatcher_2. 2021-06-02 17:13:13.516+0300 info [DefaultJobLeaderService] Stop job leader service. 2021-06-02 17:13:13.517+0300 info [TaskExecutorLocalStateStoresManager] Shutting down TaskExecutorLocalStateStoresManager. 2021-06-02 17:13:13.520+0300 info [FileChannelManagerImpl] FileChannelManager removed spill file directory /var/folders/ms/18j3nxjd6rzcj9mmgly3n0rrq_g_vb/T/flink-io-dcec76d9-4e46-411d-ad09-4257383a9df6 2021-06-02 17:13:13.521+0300 info [NettyShuffleEnvironment] Shutting down the network environment and its components. 2021-06-02 17:13:13.523+0300 info [FileChannelManagerImpl] FileChannelManager removed spill file directory /var/folders/ms/18j3nxjd6rzcj9mmgly3n0rrq_g_vb/T/flink-netty-shuffle-6412980d-faa6-4a3d-a3e0-138001f34bdc 2021-06-02 17:13:13.523+0300 info [KvStateService] Shutting down the kvState service and its components. 2021-06-02 17:13:13.523+0300 info [DefaultJobLeaderService] Stop job leader service. 2021-06-02 17:13:13.525+0300 info [FileCache] removed file cache directory /var/folders/ms/18j3nxjd6rzcj9mmgly3n0rrq_g_vb/T/flink-dist-cache-763c1c92-6ae6-4e0d-8cec-2b42b11476ac 2021-06-02 17:13:13.526+0300 info [TaskExecutor] Stopped TaskExecutor akka://flink/user/rpc/taskmanager_0. 2021-06-02 17:13:13.526+0300 info [AkkaRpcService] Stopping Akka RPC service. 2021-06-02 17:13:13.562+0300 info [AkkaRpcService] Stopping Akka RPC service. 2021-06-02 17:13:13.562+0300 info [AkkaRpcService] Stopped Akka RPC service. 2021-06-02 17:13:13.567+0300 info [PermanentBlobCache] Shutting down BLOB cache 2021-06-02 17:13:13.569+0300 info [TransientBlobCache] Shutting down BLOB cache 2021-06-02 17:13:13.574+0300 info [BlobServer] Stopped BLOB server at 0.0.0.0:56749 2021-06-02 17:13:13.574+0300 info [AkkaRpcService] Stopped Akka RPC service. Process finished with exit code 0 -- Sent from: http://apache-flink-user-mailing-list-archive.2336050.n4.nabble.com/ |
I fear that you hit FLINK-20888 [1]. I'll prioritize a fix. In the meantime, you can try to add a forward()/shuffle() write after the source. On Wed, Jun 2, 2021 at 4:30 PM Taras Moisiuk <[hidden email]> wrote: Hi Fabian, |
Thank you,
looks like shuffle() works -- Sent from: http://apache-flink-user-mailing-list-archive.2336050.n4.nabble.com/ |
Free forum by Nabble | Edit this page |