André Luiz Diniz da Silva
06/14/2023, 5:48 PMCaused 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.André Luiz Diniz da Silva
06/14/2023, 5:48 PMCaused 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) ~[?:?]André Luiz Diniz da Silva
06/14/2023, 5:49 PMval 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"
)Martijn Visser
06/14/2023, 6:06 PMAndré Luiz Diniz da Silva
06/14/2023, 6:19 PMAndré Luiz Diniz da Silva
06/14/2023, 6:20 PM2.10.1 to the classpathAndré Luiz Diniz da Silva
06/14/2023, 6:25 PM"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.Martijn Visser
06/14/2023, 6:26 PMAndré Luiz Diniz da Silva
06/14/2023, 6:31 PM+-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.4André Luiz Diniz da Silva
06/14/2023, 6:32 PMMartijn Visser
06/14/2023, 6:32 PMAndré Luiz Diniz da Silva
06/14/2023, 6:33 PMflink-avro-confluent-registryAndré Luiz Diniz da Silva
06/14/2023, 6:33 PMMartijn Visser
06/14/2023, 6:33 PMAndré Luiz Diniz da Silva
06/14/2023, 6:39 PM/hadoop/hadoop-2.10.1/dist/share/hadoop/common/lib/guava-11.0.2.jarAndré Luiz Diniz da Silva
06/14/2023, 6:40 PM