【问题标题】:Iceberg is not working when writing AVRO from spark从 spark 编写 AVRO 时 Iceberg 不工作
【发布时间】:2021-05-02 01:41:59
【问题描述】:
  1. 在将 AVRO 文件从 GCS 附加到表时遇到以下错误。 avro 文件是有效的,但我们使用的是放气的 avro,这是一个问题吗?

org.apache.iceberg.avro.AvroIterable.newFileReader(AvroIterable.java:101) 的线程“streaming-job-executor-0”java.lang.NoClassDefFoundError: org/apache/avro/InvalidAvroMagicException 中的异常。 apache.iceberg.avro.AvroIterable.iterator(AvroIterable.java:77) 在 org.apache.iceberg.avro.AvroIterable.iterator(AvroIterable.java:37) 在 org.apache.iceberg.relocated.com.google.common。 collect.Iterables.addAll(Iterables.java:320) at org.apache.iceberg.relocated.com.google.common.collect.Lists.newLinkedList(Lists.java:237) at org.apache.iceberg.ManifestLists.read( ManifestLists.java:46) at org.apache.iceberg.BaseSnapshot.cacheManifests(BaseSnapshot.java:127) at org.apache.iceberg.BaseSnapshot.dataManifests(BaseSnapshot.java:149) at org.apache.iceberg.MergingSnapshotProducer.apply (MergingSnapshotProducer.java:343) 在 org.apache.iceberg.SnapshotProducer.apply(SnapshotProducer.java:163) 在 org.apache.iceberg.SnapshotProducer.lambda$commit$2(SnapshotProducer.java:276) 在 org.apache.icebe rg.util.Tasks$Builder.runTaskWithRetry(Tasks.java:404) 在 org.apache.iceberg.util.Tasks$Builder.runSingleThreaded(Tasks.java:213) 在 org.apache.iceberg.util.Tasks$Builder。 run(Tasks.java:197) at org.apache.iceberg.util.Tasks$Builder.run(Tasks.java:189) at org.apache.iceberg.SnapshotProducer.commit(SnapshotProducer.java:275) at com.snapchat .transformer.TransformerStreamingWorker.lambda$execute$d121240d$1(TransformerStreamingWorker.java:162) 在 org.apache.spark.streaming.api.java.JavaDStreamLike.$anonfun$foreachRDD$2(JavaDStreamLike.scala:280) 在 org.apache。 spark.streaming.api.java.JavaDStreamLike.$anonfun$foreachRDD$2$adapted(JavaDStreamLike.scala:280) at org.apache.spark.streaming.dstream.ForEachDStream.$anonfun$generateJob$2(ForEachDStream.scala:51) at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) 在 org.apache.spark.streaming.dstream.DStream.createRDDWithLocalProperties(DStream.scala:416) 在 org.apache。 spark.streaming.dstream.ForEachDStream.$anonfun$generateJ ob$1(ForEachDStream.scala:51) at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) at scala.util.Try$.apply(Try.scala:213)在 org.apache.spark.streaming.scheduler.Job.run(Job.scala:39) 在 org.apache.spark.streaming.scheduler.JobScheduler$JobHandler.$anonfun$run$1(JobScheduler.scala:257) 在 scala .runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) 在 scala.util.DynamicVariable.withValue(DynamicVariable.scala:62) 在 org.apache.spark.streaming.scheduler.JobScheduler $JobHandler.run(JobScheduler.scala:257) 在 java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) 在 java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) 在 java. lang.Thread.run(Thread.java:748) 由:java.lang.ClassLoader 的 java.net.URLClassLoader.findClass(URLClassLoader.java:382) 处的 java.lang.ClassNotFoundException: org.apache.avro.InvalidAvroMagicException 引起。 loadClass(ClassLoader.java:418) 在 java.lang.ClassLoader.loadClass(ClassLoader. java:351) ... 38 更多

  1. 日志显示冰山表已经存在,但是我在gcs中看不到元数据文件?我正在从 dataproc 集群运行 spark 作业,我在哪里可以看到元数据文件?

##################### 冰山版本:0.11 火花版本 3.0 ######################

public void appendData(List<FileMetadata> publishedFiles, Schema icebergSchema) {
    TableIdentifier tableIdentifier = TableIdentifier.of(TRANSFORMER, jobConfig.streamName());
    // PartitionSpec partitionSpec = IcebergInternalFields.getPartitionSpec(tableSchema);
    HadoopTables tables = new HadoopTables(new Configuration());

   
    PartitionSpec partitionSpec = PartitionSpec.builderFor(icebergSchema)
            .build();

    Table table = null;
    if (tables.exists(tableIdentifier.name())) {
        table = tables.load(tableIdentifier.name());
    } else {
        table = tables.create(
                icebergSchema,
                partitionSpec,
                tableIdentifier.name());
    }
    AppendFiles appendFiles = table.newAppend();
    for (FileMetadata fileMetadata : publishedFiles) {

        appendFiles.appendFile(DataFiles.builder(partitionSpec)
                .withPath(fileMetadata.getFilename())
                .withFileSizeInBytes(fileMetadata.getFileSize())
                .withRecordCount(fileMetadata.getRowCount())
                .withFormat(FileFormat.AVRO)
                .build());
    }
    appendFiles.commit();
}

【问题讨论】:

    标签: apache-spark google-cloud-storage spark-avro iceberg


    【解决方案1】:

    以下两件事解决了我的问题

    • 确保我为冰山表提供了正确的路径名(在我的例子中带有 gs:// 前缀)

    • 解决了apache.avro依赖版本冲突

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2014-01-03
      • 1970-01-01
      • 2022-07-06
      • 2022-10-17
      • 2021-07-01
      相关资源
      最近更新 更多