批量写入数据(Bulkload)
Bulkload 是 Java SDK 的高吞吐批量写入接口,适合把本地文件、外部系统查询结果或历史数据批量导入 Lakehouse 表。
概述
Bulkload 的写入流程是:Java 应用创建 Row,SDK 将数据暂存到对象存储或 Table Volume,应用显式执行
getCommitRequests()
getCommitRequests()
、
prepareCommit()
prepareCommit()
、
commit()
commit()
和
future.get()
future.get()
后,数据才会加载到目标表;最后调用
close()
close()
释放资源。它适合分钟级批任务,不适合毫秒级实时写入。
| 写入方式 | 适用场景 | 数据可见性 |
|---|
RealtimeStream
RealtimeStream | 高频小批量、事件流、CDC | 写入后可较快查询,提交后下游对象可见 |
| Bulkload | 大批量导入、历史数据回灌、文件解析写入 | 显式 commit()
commit() 成功后可见 |
| JDBC DML | 低频小批量 DML | SQL 执行成功后可见 |
使用限制
- Bulkload 不适合 5 分钟以内高频触发的小批量写入,实时场景使用
RealtimeStream
RealtimeStream
。
- Bulkload 写入的数据需要显式提交成功后才可见,只调用
close()
close()
不会提交数据。
- 并发写入时,不同线程或进程应使用不同分片 ID,避免互相覆盖。
- 批量写入统一使用本文介绍的 Bulkload 接口。
创建批量写入流
import com.clickzetta.client.BulkloadStreamV2;
import com.clickzetta.client.ClickZettaClient;
import com.clickzetta.platform.bulkload.adapt.Committable;
import com.clickzetta.platform.bulkload.v2.BulkLoadCommitter;
import com.clickzetta.platform.client.api.BulkLoadOperation;
import com.clickzetta.platform.client.api.Row;
import java.util.Arrays;
import java.util.Collection;
public class BulkloadDemo {
public static void main(String[] args) throws Exception {
String jdbcUrl = args[0];
String username = args[1];
String password = args[2];
String schema = args[3];
String table = args[4];
ClickZettaClient client = ClickZettaClient.newBuilder()
.url(jdbcUrl)
.username(username)
.password(password)
.build();
BulkloadStreamV2 stream = null;
try {
stream = client.newBulkloadStreamV2Builder()
.withOperation(BulkLoadOperation.APPEND)
.createStream(client, schema, table);
Row row = stream.createRow(0);
row.setValue("id", 1);
row.setValue("name", "bulkload_001");
stream.apply(row, 0);
// 提交本次写入:收集提交请求 -> 预提交 -> 提交 -> 等待完成。
Collection<BulkLoadCommitter.CommitRequest<Committable>> commitRequests =
stream.getCommitRequests();
String transactionId = stream.prepareCommit(commitRequests);
stream.commit(Arrays.asList(transactionId), commitRequests).get();
} finally {
if (stream != null) {
stream.close();
}
client.close();
}
}
}
⚠️ 注意:Bulkload 数据必须显式提交后才会写入目标表。只调用
close()
close()
只会释放资源,不会提交数据。提交流程固定为
getCommitRequests()
getCommitRequests()
→
prepareCommit()
prepareCommit()
→
commit()
commit()
→
future.get()
future.get()
,以
future.get()
future.get()
等待提交完成。
写入分区表
静态分区写入时,通过
withPartitionSpecs
withPartitionSpecs
指定目标分区:
import com.clickzetta.client.BulkloadStreamV2;
import com.clickzetta.client.ClickZettaClient;
import com.clickzetta.platform.client.api.BulkLoadOperation;
public class BulkloadPartitionDemo {
public static void main(String[] args) throws Exception {
ClickZettaClient client = ClickZettaClient.newBuilder()
.url(args[0])
.username(args[1])
.password(args[2])
.build();
BulkloadStreamV2 stream = null;
try {
stream = client.newBulkloadStreamV2Builder()
.withOperation(BulkLoadOperation.APPEND)
.withPartitionSpecs("dt=2026-07-14,region=cn")
.createStream(client, args[3], args[4]);
} finally {
if (stream != null) {
stream.close();
}
client.close();
}
}
}
动态分区写入时,不指定
withPartitionSpecs
withPartitionSpecs
,并在 Row 中写入分区列的实际值。
⚠️ 注意:静态分区会把本次写入数据全部导入指定分区;如果数据中包含分区列的其他值,也不会自动分流到其他分区。
写入复杂类型
import com.clickzetta.client.BulkloadStreamV2;
import com.clickzetta.platform.client.api.Row;
import java.util.Arrays;
import java.util.HashMap;
import java.util.Map;
public class BulkloadComplexTypeDemo {
public static void writeRow(BulkloadStreamV2 stream) throws Exception {
Row row = stream.createRow(0);
row.setValue("tags", Arrays.asList("new", "paid"));
Map<Integer, String> properties = new HashMap<>();
properties.put(1, "mobile");
row.setValue("properties", properties);
Map<String, Object> profile = new HashMap<>();
profile.put("level", 3);
profile.put("city", "Shanghai");
row.setValue("profile", profile);
stream.apply(row, 0);
}
}
并发写入
当单次导入数据量较大时,可以用多个线程并发写入。每个并发任务使用独立分片 ID:
import com.clickzetta.client.BulkloadStreamV2;
import com.clickzetta.platform.client.api.Row;
public class BulkloadConcurrentWriteDemo {
public static void writeRow(BulkloadStreamV2 stream, int shardId) throws Exception {
Row row = stream.createRow(shardId);
row.setValue("id", 10003);
row.setValue("name", "parallel_003");
stream.apply(row, shardId);
}
}
所有并发任务写入结束后,统一执行
getCommitRequests()
getCommitRequests()
→
prepareCommit()
prepareCommit()
→
commit()
commit()
→
future.get()
future.get()
提交导入任务,最后调用
stream.close()
stream.close()
释放资源。
⚠️ 注意:不要让多个并发任务复用同一个分片 ID,否则可能导致分片数据互相覆盖。
操作类型
| 操作 | 说明 |
|---|
BulkLoadOperation.APPEND
BulkLoadOperation.APPEND | 追加写入目标表 |
BulkLoadOperation.OVERWRITE
BulkLoadOperation.OVERWRITE | 覆盖目标表或目标分区 |
BulkLoadOperation.UPSERT
BulkLoadOperation.UPSERT | 按 recordKeys
recordKeys 指定的记录键进行插入更新;使用 UPSERT
UPSERT 时必须通过 withRecordKey
withRecordKey 或 withRecordKeys
withRecordKeys 指定记录键,partialUpdateColumns
partialUpdateColumns 仅支持在 UPSERT
UPSERT 模式下配置 |
相关文档