You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.04 06:55:16