apache / apache/hudi

flink offlinecompact error

Open
#9,500 3 comments 0 reactions 0 assignees View on GitHub
area:table-service engine:flink priority:high
Dominant language
Java
Stars
6.2k
Forks
2.5k
Avg merge
2d 8h
Merged PRs (30d)
111

Description

Hello, I encountered an error when I use flink offline compact :

2023-08-21 16:03:22,756 INFO org.apache.flink.runtime.taskmanager.Task [] - compact_task (3/3)#0 (fd49759c8ac941659bb460ca82108333) switched from INITIALIZING to RUNNING.
2023-08-21 16:03:22,757 INFO org.apache.flink.runtime.taskmanager.Task [] - compact_task (2/3)#0 (6caf7347aa5869e2ab11b3d628871ac9) switched from INITIALIZING to RUNNING.
2023-08-21 16:03:23,896 WARN org.apache.flink.runtime.taskmanager.Task [] - compact_task (2/3)#0 (6caf7347aa5869e2ab11b3d628871ac9) switched from RUNNING to FAILED with failure cause: java.lang.RuntimeException: Cannot instantiate class.
at org.apache.flink.api.java.typeutils.runtime.PojoSerializer.createInstance(PojoSerializer.java:213)
at org.apache.flink.api.java.typeutils.runtime.PojoSerializer.deserialize(PojoSerializer.java:413)
at org.apache.flink.streaming.runtime.streamrecord.StreamElementSerializer.deserialize(StreamElementSerializer.java:193)
at org.apache.flink.streaming.runtime.streamrecord.StreamElementSerializer.deserialize(StreamElementSerializer.java:46)
at org.apache.flink.runtime.plugable.NonReusingDeserializationDelegate.read(NonReusingDeserializationDelegate.java:53)
at org.apache.flink.runtime.io.network.api.serialization.NonSpanningWrapper.readInto(NonSpanningWrapper.java:337)
at org.apache.flink.runtime.io.network.api.serialization.SpillingAdaptiveSpanningRecordDeserializer.readNonSpanningRecord(SpillingAdaptiveSpanningRecordDeserializer.java:128)
at org.apache.flink.runtime.io.network.api.serialization.SpillingAdaptiveSpanningRecordDeserializer.readNextRecord(SpillingAdaptiveSpanningRecordDeserializer.java:103)
at org.apache.flink.runtime.io.network.api.serialization.SpillingAdaptiveSpanningRecordDeserializer.getNextRecord(SpillingAdaptiveSpanningRecordDeserializer.java:93)
at org.apache.flink.streaming.runtime.io.AbstractStreamTaskNetworkInput.emitNext(AbstractStreamTaskNetworkInput.java:95)
at org.apache.flink.streaming.runtime.io.StreamOneInputProcessor.processInput(StreamOneInputProcessor.java:65)
at org.apache.flink.streaming.runtime.tasks.StreamTask.processInput(StreamTask.java:496)
at org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.runMailboxLoop(MailboxProcessor.java:203)
at org.apache.flink.streaming.runtime.tasks.StreamTask.runMailboxLoop(StreamTask.java:809)
at org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:761)
at org.apache.flink.runtime.taskmanager.Task.runWithSystemExitMonitoring(Task.java:958)
at org.apache.flink.runtime.taskmanager.Task.restoreAndInvoke(Task.java:937)
at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:766)
at org.apache.flink.runtime.taskmanager.Task.run(Task.java:575)
at java.lang.Thread.run(Thread.java:748)
Caused by: java.lang.IllegalArgumentException: Can not set org.apache.hudi.common.model.CompactionOperation field org.apache.hudi.sink.compact.CompactionPlanEvent.operation to org.apache.hudi.common.model.CompactionOperation
at sun.reflect.UnsafeFieldAccessorImpl.throwSetIllegalArgumentException(UnsafeFieldAccessorImpl.java:167)
at sun.reflect.UnsafeFieldAccessorImpl.throwSetIllegalArgumentException(UnsafeFieldAccessorImpl.java:171)
at sun.reflect.UnsafeObjectFieldAccessorImpl.set(UnsafeObjectFieldAccessorImpl.java:81)
at java.lang.reflect.Field.set(Field.java:764)
at org.apache.flink.api.java.typeutils.runtime.PojoSerializer.initializeFields(PojoSerializer.java:221)
at org.apache.flink.api.java.typeutils.runtime.PojoSerializer.createInstance(PojoSerializer.java:210)
... 19 more

./bin/flink run \
-c org.apache.hudi.sink.compact.HoodieFlinkCompactor \
lib/hudi-flink1.14-bundle-0.12.0.jar \
--compaction-max-memory 1024 \
--compaction-tasks 3 \
--path hdfs:/.../tableA

* Hudi version : 0.12.0
* Flink version : 0.14.5
* Scala version: 2.12
* Hadoop version : 3.2.1
* Run mode : flink on yarn session

when I run it on my local VM centos, it's ok ; but run product env, it appears error,
please Help me !

Contributor guide

No contributing guide indexed for this repository

Research direction

Start at org.apache.hudi.sink.compact.HoodieFlinkCompactor and inspect the compaction event deserialization shown in the stack trace, especially CompactionPlanEvent.operation and Flink's PojoSerializer. Compare the local VM and production Flink-on-YARN environments using the reported Hudi 0.12.0 bundle and dependency versions. Done means identifying and reproducing the environment-specific failure, then verifying that offline compaction completes successfully.

Written by the indexing model from the issue text.

Assessment

Tech stack
java
Domain
data-engineering, stream-processing
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Stale
Clarity
Needs clarification
Newbie friendliness
25/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.