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

求助:如何用Delta Lake Standalone创建Delta表并写入数据?求Scala示例

使用Delta Lake Standalone API写入S3 Delta表(Scala示例)

Delta Standalone API本身不提供数据写入能力,核心流程是:先将数据写入Parquet格式文件,再通过Delta Standalone提交元数据到Delta日志,完成Delta表的创建/更新。以下是完整的Scala实现示例,包含Parquet数据写入、AddFile元数据构建、Delta日志提交全流程。

补充依赖

除你已有的依赖外,需添加Parquet写入及S3适配相关依赖:

<dependency>
    <groupId>org.apache.parquet</groupId>
    <artifactId>parquet-hadoop</artifactId>
    <version>1.12.3</version>
</dependency>
<dependency>
    <groupId>org.apache.hadoop</groupId>
    <artifactId>hadoop-aws</artifactId>
    <version>3.3.1</version>
</dependency>

完整Scala代码示例

import org.apache.hadoop.conf.Configuration
import org.apache.hadoop.fs.{FileSystem, Path}
import org.apache.parquet.hadoop.ParquetWriter
import org.apache.parquet.hadoop.api.WriteSupport
import org.apache.parquet.hadoop.metadata.CompressionCodecName
import org.apache.parquet.schema.MessageTypeParser
import io.delta.standalone.DeltaLog
import io.delta.standalone.actions.{AddFile, CommitInfo}
import io.delta.standalone.operations.{WriteBuilder, WriteOperation}

import java.util.Collections

object DeltaStandaloneWriterExample {
  def main(args: Array[String]): Unit = {
    // 1. 初始化Hadoop配置(适配S3)
    val conf = new Configuration()
    // 替换为你的S3实际配置
    conf.set("fs.s3a.access.key", "your-s3-access-key")
    conf.set("fs.s3a.secret.key", "your-s3-secret-key")
    conf.set("fs.s3a.endpoint", "s3.amazonaws.com")
    conf.set("fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem")

    // Delta表在S3的根路径
    val deltaTableRoot = "s3a://your-bucket/target-delta-table"
    // 待写入的Parquet文件路径(建议放在Delta表目录下的data子目录)
    val parquetFileRelPath = "data/part-00000.parquet"
    val parquetFullPath = new Path(s"$deltaTableRoot/$parquetFileRelPath")

    // 2. 定义Parquet Schema(示例结构:id: int, name: string, value: double)
    val schemaStr = """message sample_data {
                      |  required int32 id;
                      |  required binary name (UTF8);
                      |  required double value;
                      |}""".stripMargin
    val parquetSchema = MessageTypeParser.parseMessageType(schemaStr)

    // 3. 写入Parquet数据到S3
    val writeSupport = WriteSupport.forSchema(parquetSchema, conf)
    val writer = new ParquetWriter.Builder[Array[Any]](parquetFullPath, writeSupport)
      .withConf(conf)
      .withCompressionCodec(CompressionCodecName.SNAPPY)
      .build()

    // 写入示例数据
    writer.write(Array(1, "Alice", 100.5))
    writer.write(Array(2, "Bob", 200.7))
    writer.close()

    // 4. 收集AddFile所需元数据
    val fs = FileSystem.get(parquetFullPath.toUri, conf)
    val fileStatus = fs.getFileStatus(parquetFullPath)
    val fileSize = fileStatus.getLen

    // 读取Parquet统计信息(可选,用于Delta查询优化)
    val parquetMeta = org.apache.parquet.hadoop.ParquetFileReader.readFooter(conf, parquetFullPath)
    val blockStats = parquetMeta.getBlocks.get(0).getStatistics

    // 构建AddFile对象
    val addFile = AddFile.builder()
      .path(parquetFileRelPath) // 必须是相对于Delta表根目录的路径
      .partitionValues(Collections.emptyMap()) // 无分区则传空Map,有分区则传入键值对
      .size(fileSize)
      .modificationTime(System.currentTimeMillis())
      .dataChange(true)
      // 手动构造统计信息JSON,也可通过Parquet元数据自动生成
      .stats(s"""{"numRecords":2,"minValues":{"id":1,"name":"Alice","value":100.5},"maxValues":{"id":2,"name":"Bob","value":200.7}}""")
      .build()

    // 5. 提交元数据到Delta日志
    val deltaLog = DeltaLog.forTable(conf, deltaTableRoot)
    val writeBuilder: WriteBuilder = deltaLog.startTransaction()
      .newWriteBuilder()
      .operation(WriteOperation.WRITE)
      .addActions(Collections.singletonList(addFile)) // 支持批量添加多个AddFile
      .commitInfo(CommitInfo.builder().userId("standalone-writer").build())

    // 执行事务提交
    writeBuilder.commit()

    println("Delta表创建/更新完成")
  }
}

关键细节说明

  • AddFile路径规则:必须使用相对于Delta表根目录的相对路径,不能用绝对路径,否则Delta日志无法正确关联文件。
  • 统计信息(stats):虽为可选字段,但添加后能让Delta Lake跳过不必要的文件扫描,大幅提升查询性能,建议通过Parquet元数据自动生成或手动构造符合格式的JSON。
  • S3配置校验:确保Hadoop配置中的S3密钥、endpoint等参数正确,否则会出现权限拒绝或路径无法访问的错误。
  • 事务性保障:startTransaction()会获取当前日志版本,提交时自动处理版本冲突,保证Delta表的ACID特性。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 13:18:20