求助:如何用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
相关产品推荐
相关产品推荐

