Hello everyone! how are you? I’m having a problem ...
# troubleshooting
a
Hello everyone! how are you? I’m having a problem related I think to class loading. I ’m getting the following error:
Caused by: java.lang.NoClassDefFoundError: com/google/common/util/concurrent/internal/InternalFutureFailureAccess
(complete stack trace in the thread). The context is that I have a job that write parquet files and after adding a new sink to write data to a Kafka topic using avro with schema registry this error began to happen. Any ideia what could it be? am I missing some dependency? More details in the thread.
the stacktrace:
Copy code
Caused by: java.lang.NoClassDefFoundError: com/google/common/util/concurrent/internal/InternalFutureFailureAccess
	at com.google.common.cache.LocalCache$LoadingValueReference.<init>(LocalCache.java:3472) ~[magneton-assembly-0.1.0-SNAPSHOT.jar:0.1.0-SNAPSHOT]
	at com.google.common.cache.LocalCache$LoadingValueReference.<init>(LocalCache.java:3476) ~[magneton-assembly-0.1.0-SNAPSHOT.jar:0.1.0-SNAPSHOT]
	at com.google.common.cache.LocalCache$Segment.lockedGetOrLoad(LocalCache.java:2134) ~[magneton-assembly-0.1.0-SNAPSHOT.jar:0.1.0-SNAPSHOT]
	at com.google.common.cache.LocalCache$Segment.get(LocalCache.java:2045) ~[magneton-assembly-0.1.0-SNAPSHOT.jar:0.1.0-SNAPSHOT]
	at com.google.common.cache.LocalCache.get(LocalCache.java:3962) ~[magneton-assembly-0.1.0-SNAPSHOT.jar:0.1.0-SNAPSHOT]
	at com.google.common.cache.LocalCache.getOrLoad(LocalCache.java:3985) ~[magneton-assembly-0.1.0-SNAPSHOT.jar:0.1.0-SNAPSHOT]
	at com.google.common.cache.LocalCache$LocalLoadingCache.get(LocalCache.java:4946) ~[magneton-assembly-0.1.0-SNAPSHOT.jar:0.1.0-SNAPSHOT]
	at com.google.common.cache.LocalCache$LocalLoadingCache.getUnchecked(LocalCache.java:4952) ~[magneton-assembly-0.1.0-SNAPSHOT.jar:0.1.0-SNAPSHOT]
	at org.apache.hadoop.io.compress.CodecPool.updateLeaseCount(CodecPool.java:135) ~[hadoop-common-2.10.1.jar:?]
	at org.apache.hadoop.io.compress.CodecPool.getCompressor(CodecPool.java:162) ~[hadoop-common-2.10.1.jar:?]
	at org.apache.hadoop.io.compress.CodecPool.getCompressor(CodecPool.java:168) ~[hadoop-common-2.10.1.jar:?]
	at org.apache.parquet.hadoop.CodecFactory$HeapBytesCompressor.<init>(CodecFactory.java:146) ~[magneton-assembly-0.1.0-SNAPSHOT.jar:0.1.0-SNAPSHOT]
	at org.apache.parquet.hadoop.CodecFactory.createCompressor(CodecFactory.java:208) ~[magneton-assembly-0.1.0-SNAPSHOT.jar:0.1.0-SNAPSHOT]
	at org.apache.parquet.hadoop.CodecFactory.getCompressor(CodecFactory.java:191) ~[magneton-assembly-0.1.0-SNAPSHOT.jar:0.1.0-SNAPSHOT]
	at org.apache.parquet.hadoop.ParquetWriter.<init>(ParquetWriter.java:296) ~[magneton-assembly-0.1.0-SNAPSHOT.jar:0.1.0-SNAPSHOT]
	at org.apache.parquet.hadoop.ParquetWriter$Builder.build(ParquetWriter.java:675) ~[magneton-assembly-0.1.0-SNAPSHOT.jar:0.1.0-SNAPSHOT]
	at org.apache.flink.formats.parquet.row.ParquetRowDataBuilder$FlinkParquetBuilder.createWriter(ParquetRowDataBuilder.java:140) ~[magneton-assembly-0.1.0-SNAPSHOT.jar:0.1.0-SNAPSHOT]
	at org.apache.flink.formats.parquet.ParquetWriterFactory.create(ParquetWriterFactory.java:56) ~[magneton-assembly-0.1.0-SNAPSHOT.jar:0.1.0-SNAPSHOT]
	at org.apache.flink.connector.file.table.FileSystemTableSink$ProjectionBulkFactory.create(FileSystemTableSink.java:637) ~[flink-connector-files-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.api.functions.sink.filesystem.BulkBucketWriter.openNew(BulkBucketWriter.java:76) ~[magneton-assembly-0.1.0-SNAPSHOT.jar:0.1.0-SNAPSHOT]
	at org.apache.flink.streaming.api.functions.sink.filesystem.OutputStreamBasedPartFileWriter$OutputStreamBasedBucketWriter.openNewInProgressFile(OutputStreamBasedPartFileWriter.java:124) ~[magneton-assembly-0.1.0-SNAPSHOT.jar:0.1.0-SNAPSHOT]
	at org.apache.flink.streaming.api.functions.sink.filesystem.BulkBucketWriter.openNewInProgressFile(BulkBucketWriter.java:36) ~[magneton-assembly-0.1.0-SNAPSHOT.jar:0.1.0-SNAPSHOT]
	at org.apache.flink.streaming.api.functions.sink.filesystem.Bucket.rollPartFile(Bucket.java:244) ~[flink-dist-1.17.1.jar:0.1.0-SNAPSHOT]
	at org.apache.flink.streaming.api.functions.sink.filesystem.Bucket.write(Bucket.java:221) ~[flink-dist-1.17.1.jar:0.1.0-SNAPSHOT]
	at org.apache.flink.streaming.api.functions.sink.filesystem.Buckets.onElement(Buckets.java:305) ~[flink-dist-1.17.1.jar:0.1.0-SNAPSHOT]
	at org.apache.flink.streaming.api.functions.sink.filesystem.StreamingFileSinkHelper.onElement(StreamingFileSinkHelper.java:103) ~[flink-dist-1.17.1.jar:0.1.0-SNAPSHOT]
	at org.apache.flink.connector.file.table.stream.AbstractStreamingWriter.processElement(AbstractStreamingWriter.java:140) ~[flink-connector-files-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.runtime.tasks.CopyingChainingOutput.pushToOperator(CopyingChainingOutput.java:75) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.runtime.tasks.CopyingChainingOutput.collect(CopyingChainingOutput.java:50) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.runtime.tasks.CopyingChainingOutput.collect(CopyingChainingOutput.java:29) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.table.runtime.operators.sink.StreamRecordTimestampInserter.processElement(StreamRecordTimestampInserter.java:52) ~[flink-table-runtime-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.runtime.tasks.CopyingChainingOutput.pushToOperator(CopyingChainingOutput.java:75) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.runtime.tasks.CopyingChainingOutput.collect(CopyingChainingOutput.java:50) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.runtime.tasks.CopyingChainingOutput.collect(CopyingChainingOutput.java:29) ~[flink-dist-1.17.1.jar:1.17.1]
	at StreamExecCalc$99.processElement_split45(Unknown Source) ~[?:?]
	at StreamExecCalc$99.processElement(Unknown Source) ~[?:?]
	at org.apache.flink.streaming.runtime.tasks.CopyingChainingOutput.pushToOperator(CopyingChainingOutput.java:75) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.runtime.tasks.CopyingChainingOutput.collect(CopyingChainingOutput.java:50) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.runtime.tasks.CopyingChainingOutput.collect(CopyingChainingOutput.java:29) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.table.runtime.operators.source.InputConversionOperator.processElement(InputConversionOperator.java:128) ~[flink-table-runtime-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.runtime.tasks.CopyingChainingOutput.pushToOperator(CopyingChainingOutput.java:75) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.runtime.tasks.CopyingChainingOutput.collect(CopyingChainingOutput.java:50) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.runtime.tasks.CopyingChainingOutput.collect(CopyingChainingOutput.java:29) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.api.operators.StreamMap.processElement(StreamMap.java:38) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.runtime.tasks.CopyingChainingOutput.pushToOperator(CopyingChainingOutput.java:75) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.runtime.tasks.CopyingChainingOutput.collect(CopyingChainingOutput.java:50) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.runtime.tasks.CopyingChainingOutput.collect(CopyingChainingOutput.java:29) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.api.operators.TimestampedCollector.collect(TimestampedCollector.java:51) ~[flink-dist-1.17.1.jar:1.17.1]
	at avalanche.core.flink.api.datastream.functions.core$package$.avalanche$core$flink$api$datastream$functions$core$package$$anon$4$$_$flatMap$$anonfun$1(core.scala:48) ~[magneton-assembly-0.1.0-SNAPSHOT.jar:0.1.0-SNAPSHOT]
	at scala.runtime.function.JProcedure1.apply(JProcedure1.java:15) ~[magneton-assembly-0.1.0-SNAPSHOT.jar:0.1.0-SNAPSHOT]
	at scala.runtime.function.JProcedure1.apply(JProcedure1.java:10) ~[magneton-assembly-0.1.0-SNAPSHOT.jar:0.1.0-SNAPSHOT]
	at scala.collection.IterableOnceOps.foreach(IterableOnce.scala:575) ~[magneton-assembly-0.1.0-SNAPSHOT.jar:0.1.0-SNAPSHOT]
	at scala.collection.IterableOnceOps.foreach$(IterableOnce.scala:573) ~[magneton-assembly-0.1.0-SNAPSHOT.jar:0.1.0-SNAPSHOT]
	at scala.collection.AbstractIterator.foreach(Iterator.scala:1300) ~[magneton-assembly-0.1.0-SNAPSHOT.jar:0.1.0-SNAPSHOT]
	at avalanche.core.flink.api.datastream.functions.core$package$$anon$4.flatMap(core.scala:48) ~[magneton-assembly-0.1.0-SNAPSHOT.jar:0.1.0-SNAPSHOT]
	at org.apache.flink.streaming.api.operators.StreamFlatMap.processElement(StreamFlatMap.java:47) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.runtime.tasks.CopyingChainingOutput.pushToOperator(CopyingChainingOutput.java:75) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.runtime.tasks.CopyingChainingOutput.collect(CopyingChainingOutput.java:50) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.runtime.tasks.CopyingChainingOutput.collect(CopyingChainingOutput.java:29) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.api.operators.TimestampedCollector.collect(TimestampedCollector.java:51) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.api.operators.async.queue.StreamRecordQueueEntry.emitResult(StreamRecordQueueEntry.java:64) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.api.operators.async.queue.OrderedStreamElementQueue.emitCompletedElement(OrderedStreamElementQueue.java:71) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.api.operators.async.AsyncWaitOperator.outputCompletedElement(AsyncWaitOperator.java:393) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.api.operators.async.AsyncWaitOperator.access$1800(AsyncWaitOperator.java:92) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.api.operators.async.AsyncWaitOperator$ResultHandler.processResults(AsyncWaitOperator.java:621) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.api.operators.async.AsyncWaitOperator$ResultHandler.lambda$processInMailbox$0(AsyncWaitOperator.java:602) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$1.runThrowing(StreamTaskActionExecutor.java:50) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.runtime.tasks.mailbox.Mail.run(Mail.java:90) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.runMail(MailboxProcessor.java:398) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.processMailsNonBlocking(MailboxProcessor.java:383) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.processMail(MailboxProcessor.java:345) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.runMailboxLoop(MailboxProcessor.java:229) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.runtime.tasks.StreamTask.runMailboxLoop(StreamTask.java:839) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:788) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.runtime.taskmanager.Task.runWithSystemExitMonitoring(Task.java:952) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.runtime.taskmanager.Task.restoreAndInvoke(Task.java:931) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:745) ~[flink-dist-1.17.1.jar:1.17.1]
	at org.apache.flink.runtime.taskmanager.Task.run(Task.java:562) ~[flink-dist-1.17.1.jar:1.17.1]
	at java.lang.Thread.run(Unknown Source) ~[?:?]
The flink dependencies that I use (Sbt as a buildtool):
Copy code
val flinkVersion = "1.17.1"
val flinkDependencies = Seq(
  "org.apache.flink" % "flink-streaming-java" % flinkVersion % "provided",
  "org.apache.flink" % "flink-clients" % flinkVersion % "provided",
  "org.apache.flink" % "flink-table-common" % flinkVersion % "provided",
  "org.apache.flink" % "flink-table-runtime" % flinkVersion % "provided",
  "org.apache.flink" % "flink-table-planner-loader" % flinkVersion % "provided",
  "org.apache.flink" % "flink-table-api-java" % flinkVersion % "provided",
  "org.apache.flink" % "flink-table-api-java-bridge" % flinkVersion % "provided",
  "org.apache.flink" % "flink-connector-kafka" % flinkVersion,
  "org.apache.flink" % "flink-parquet" % flinkVersion,
  "org.apache.flink" % "flink-json" % flinkVersion,
  "org.apache.flink" % "flink-connector-files" % flinkVersion,
  "org.apache.flink" % "flink-metrics-dropwizard" % flinkVersion,
  "org.apache.flink" % "flink-statebackend-rocksdb" % flinkVersion,
  "org.apache.flink" % "flink-avro" % flinkVersion,
  "org.apache.flink" % "flink-avro-confluent-registry" % flinkVersion,
  "io.findify" %% "flink-adt" % "0.6.1"
)
m
If you want to use Parquet, you will also need to include some Hadoop dependencies
a
Hello Martjin they are already in the classpath. Before adding this new sink the parquet writing was working fine.
we add hadoop
2.10.1
to the classpath
the only difference between the two versions of the application are these two dependencies:
Copy code
"org.apache.flink" % "flink-avro" % flinkVersion,
  "org.apache.flink" % "flink-avro-confluent-registry" % flinkVersion,
I’m trying to investigate the classpath if there is anything missing but without success.
m
I think there might be a conflict on your classpath, specifically on Guava
a
taking a look at the dependency tree of my project the only guava reference is in this new library:
Copy code
+-org.apache.flink:flink-avro-confluent-registry:1.17.1
[info]   | +-com.google.code.findbugs:jsr305:1.3.9
[info]   | +-io.confluent:kafka-schema-registry-client:7.2.2
[info]   | | +-com.fasterxml.jackson.core:jackson-databind:2.13.2.2
[info]   | | | +-com.fasterxml.jackson.core:jackson-annotations:2.13.2
[info]   | | | +-com.fasterxml.jackson.core:jackson-core:2.13.2
[info]   | | | 
[info]   | | +-com.google.guava:guava:30.1.1-jre
[info]   | | +-io.confluent:common-utils:7.2.2
[info]   | | | +-org.slf4j:slf4j-api:1.7.36
[info]   | | | 
[info]   | | +-org.apache.commons:commons-compress:1.21
[info]   | | +-org.apache.kafka:kafka-clients:7.2.2-ccs
[info]   | |   +-com.github.luben:zstd-jni:1.5.2-1
[info]   | |   +-org.lz4:lz4-java:1.8.0
[info]   | |   +-org.slf4j:slf4j-api:1.7.36
[info]   | |   +-org.xerial.snappy:snappy-java:1.1.8.4
the version that is being used is 30.1.1.
m
Yeah and I think there might be a different Guava used by one of the other dependencies
a
hmmm interesting, in my project the only thing that uses guava is the
flink-avro-confluent-registry
and we only add Hadoop to the job classpath so, in theory some lib inside hadoop ecosystem is messing it up.
m
I can't check, but I would expect that either Parquet or Hadoop also has Guava in there
a
taking a look on hadoop jar files I have the following one:
Copy code
/hadoop/hadoop-2.10.1/dist/share/hadoop/common/lib/guava-11.0.2.jar
that might be the problem I will take a look.
👍 1