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

如何不借助AWS额外服务用Apache Spark将Apache Iceberg数据导入AWS Neptune?是否可扩展?

Iceberg表数据通过Spark加载到AWS Neptune的实现方案

需求背景

我正在推进一个项目,需要将Apache Iceberg表中的数据加载到AWS Neptune图数据库中。要求使用Apache Spark完成该任务,且不借助AWS DMS等额外AWS服务。项目涉及的数据量初始为400GB,之后每日约有15GB增量数据,需加载包含顶点和边数据的多张表。

环境配置与需求概述

  • 数据源:Apache Iceberg表
  • 目标数据库:AWS Neptune
  • 中间处理工具:Apache Spark

现有实现方案

我已经编写了Scala版本的Spark程序,读取Apache Iceberg表数据并通过Gremlin API将数据插入Neptune,简化后的代码如下:

import org.apache.spark.sql.{SparkSession, Row, DataFrame}
import org.apache.spark.sql.functions._
import org.apache.spark.sql.types.{StringType, StructField, StructType}
import org.apache.iceberg.spark.Spark3Util
import org.apache.tinkerpop.gremlin.driver.{Client, Cluster}
import org.apache.tinkerpop.gremlin.driver.remote.DriverRemoteConnection
import org.apache.tinkerpop.gremlin.process.remote.RemoteGraph
import org.apache.tinkerpop.gremlin.structure.Graph
import software.amazon.awssdk.auth.credentials.{AwsBasicCredentials, StaticCredentialsProvider}
import software.amazon.awssdk.regions.Region
import software.amazon.awssdk.services.neptune.NeptuneClient

object NeptuneSparkApp {

  def main(args: Array[String]): Unit = {
    // 初始化Spark会话
    val spark = SparkSession.builder
      .appName("NeptuneSparkApp")
      .config("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions")
      .config("spark.sql.catalog.spark_catalog", "org.apache.iceberg.spark.SparkSessionCatalog")
      .config("spark.sql.catalog.spark_catalog.type", "hive")
      .config("spark.sql.catalog.spark_catalog.uri", "thrift://localhost:9083") // 更新为你的Hive元存储URI
      .getOrCreate()

    // 定义AWS Neptune端点与凭证
    val neptuneEndpoint = "your-neptune-endpoint:8182"
    val neptuneRegion = "your-region"
    val awsAccessKeyId = "your-access-key-id"
    val awsSecretAccessKey = "your-secret-access-key"

    // 初始化Neptune客户端
    val credentials = AwsBasicCredentials.create(awsAccessKeyId, awsSecretAccessKey)
    val neptuneClient = NeptuneClient.builder()
      .region(Region.of(neptuneRegion))
      .credentialsProvider(StaticCredentialsProvider.create(credentials))
      .build()

    // 连接到Neptune
    val cluster = Cluster.build(neptuneEndpoint).create()
    val client = cluster.connect()
    val graph = RemoteGraph.open(client, "g")

    // 读取Iceberg表数据
    val icebergTable: DataFrame = spark.read
      .format("iceberg")
      .load("spark_catalog.db.table_name") // 更新为你的Iceberg表路径

    // 向Neptune插入数据的函数
    def insertDataToNeptune(row: Row): Unit = {
      val g = graph.traversal()
      g.addV("person")
        .property("id", row.getAs[String]("id"))
        .property("name", row.getAs[String]("name"))
        .property("age", row.getAs[Int]("age"))
        .next()
    }

    // 对DataFrame中的每一行应用插入函数
    icebergTable.collect().foreach(insertDataToNeptune)

    // 关闭连接
    client.close()
    cluster.close()

    // 停止Spark会话
    spark.stop()
  }
}

内容的提问来源于stack exchange,提问作者John Thomas

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 21:27:01