使用Java自定义Iceberg Catalog插入数据时遇ClassNotFoundException问题
自定义Iceberg Catalog插入数据时出现ClassNotFoundException
我们使用Java自定义Iceberg Catalog,表可正常创建,但插入数据时出现异常。本地Java主程序运行代码完全正常,但打包成Jar后通过Spark调用时抛出如下异常。
相关代码
public String createCustomTable(String tableName) { try { TableIdentifier tableIdentifier = TableIdentifier.of(name(), tableName); Schema schema = readSchema(tableIdentifier); Map<String, String> properties = ImmutableMap.of( TableProperties.DEFAULT_FILE_FORMAT, FileFormat.PARQUET.name() ); PartitionSpec partitionSpec = PartitionSpec.builderFor(schema) .identity(getPartitionKeyfromSchema(tableIdentifier.name())) .build(); String tableLocation = defaultLocation + tableIdentifier.namespace().toString() + "/" + tableIdentifier.name(); catalog.createTable(tableIdentifier, schema, partitionSpec, tableLocation, properties); catalog.loadTable(TableIdentifier.of(name(), tableName)); return "Table created"; } catch (Exception e) { return e.getMessage(); } } public String insertData(String tableName, String csvPath) throws IOException { Table icebergTable = catalog.loadTable(TableIdentifier.of(name(), tableName)); SparkSession spark = SparkSession.builder() .config("spark.master", "local") .getOrCreate(); String headerJson = readHeaderJson(tableName); LOGGER.info("Header JSON for {}: {}", tableName, headerJson); String[] columns = headerJson.split(","); Dataset<Row> df = spark.read() .option("header", "false") .option("inferSchema", "false") .option("comment", "#") .option("sep", "|") .csv(csvPath) .toDF(columns); LOGGER.info("Actual columns: {}", Arrays.toString(df.columns())); for (String col : df.columns()) { df = df.withColumn(col, df.col(col).cast("string")); } df.write().format("iceberg").mode(SaveMode.Append).save(icebergTable.location()); LOGGER.info("Data inserted successfully into table: {}", tableName); }
异常信息
ERROR:root:Error: An error occurred while calling o0.insertData. : java.lang.ClassNotFoundException: Failed to find data source: iceberg. Please find packages at http://spark.apache.org/third-party-projects.html at org.apache.spark.sql.errors.QueryExecutionErrors$.failedToFindDataSourceError(QueryExecutionErrors.scala:443) at org.apache.spark.sql.execution.datasources.DataSource$.lookupDataSource(DataSource.scala:670) at org.apache.spark.sql.execution.datasources.DataSource$.lookupDataSourceV2(DataSource.scala:720) at org.apache.spark.sql.DataFrameWriter.lookupV2Provider(DataFrameWriter.scala:852) at org.apache.spark.sql.DataFrameWriter.saveInternal(DataFrameWriter.scala:256) at org.apache.spark.sql.DataFrameWriter.save(DataFrameWriter.scala:239) at com.xyz.catalog.CustomCatalog.insertData(CustomCatalog.java:178) at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.base/java.lang.reflect.Method.invoke(Method.java:566) at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244) at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357) at py4j.Gateway.invoke(Gateway.java:282) at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132) at py4j.commands.CallCommand.execute(CallCommand.java:79) at py4j.GatewayConnection.run(GatewayConnection.java:238) at java.base/java.lang.Thread.run(Thread.java:829) Caused by: java.lang.ClassNotFoundException: iceberg.DefaultSource at java.base/jdk.internal.loader.BuiltinClassLoader.loadClass(BuiltinClassLoader.java:581) at java.base/jdk.internal.loader.ClassLoaders$AppClassLoader.loadClass(ClassLoaders.java:178) at java.base/java.lang.ClassLoader.loadClass(ClassLoader.java:527) at org.apache.spark.sql.execution.datasources.DataSource$.$anonfun$lookupDataSource$5(DataSource.scala:656) at scala.util.Try$.apply(Try.scala:210) at org.apache.spark.sql.execution.datasources.DataSource$.$anonfun$lookupDataSource$4(DataSource.scala:656) at scala.util.Failure.orElse(Try.scala:221) at org.apache.spark.sql.execution.datasources.DataSource$.lookupDataSource(DataSource.scala:656) ... 16 more
期望
通过df.write().format("iceberg").mode(SaveMode.Append).save(icebergTable.location());语句正常完成数据插入。
内容的提问来源于stack exchange,提问作者Mahendra Chandrasekhar
相关产品推荐
相关产品推荐

