Parquet reading/writing

classic Classic list List threaded Threaded
6 messages Options
Reply | Threaded
Open this post in threaded view
|

Parquet reading/writing

Taras Moisiuk
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/
Reply | Threaded
Open this post in threaded view
|

Re: Parquet reading/writing

Taras Moisiuk
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/
Reply | Threaded
Open this post in threaded view
|

Re: Parquet reading/writing

Fabian Paul
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/

Reply | Threaded
Open this post in threaded view
|

Re: Parquet reading/writing

Taras Moisiuk
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/

Reply | Threaded
Open this post in threaded view
|

Re: Parquet reading/writing

Arvid Heise-4
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 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/

Reply | Threaded
Open this post in threaded view
|

Re: Parquet reading/writing

Taras Moisiuk
Thank you,

looks like shuffle() works



--
Sent from: http://apache-flink-user-mailing-list-archive.2336050.n4.nabble.com/