使用 Java SDK 批量上传数据
本文档主要介绍如何使用 Java SDK 的 Bulkload 批量将数据加载到 Lakehouse 中。它适合一次性大量数据导入,支持自定义数据源,提供了数据导入的灵活性。本次案例以读取本地文件为例。如果你的数据源在对象存储或 Lakehouse Studio 数据集成支持的范围内,推荐使用 COPY 命令或数据集成功能。
参考文档
Java SDK 批量写入数据
应用场景
- 适用于需要批量上传大量数据的业务场景。
- 适合熟悉 Java 并需要自定义数据导入逻辑的开发人员。
使用限制
- BulkloadStream 不适合主键(PK)表的高频 CDC 写入;这类场景建议使用
RealtimeStream
RealtimeStream
CDC。
- 不适用于时间间隔小于五分钟的频繁数据上传场景。
使用案例
本案例以读取本地 CSV 文件为例,使用的数据集是 巴西电子商务 公共数据集中的 olist_order_items_dataset 数据。若数据源位于对象存储或 Lakehouse Studio 数据集成功能支持的范围,推荐使用 COPY 命令或数据集成功能。
前置条件
-
创建表
-
CREATE TABLE bulk_order_items (
order_id STRING,
order_item_id INT,
product_id STRING,
seller_id STRING,
shipping_limit_date STRING,
price DOUBLE,
freight_value DOUBLE
);
-
对目标表具有
INSERT
INSERT
权限。
使用 Java 代码开发
Maven依赖
在项目的
pom.xml
pom.xml
文件中添加 Lakehouse 的 Maven 依赖。Lakehouse Maven 依赖的最新版本可以在
maven库 中找到。
<dependency>
<groupId>com.clickzetta</groupId>
<artifactId>clickzetta-java</artifactId>
<version>${clickzetta-java.version}</version>
</dependency>
编写Java代码
- 初始化 Lakehouse 客户端和 BulkloadStreamV2:创建
BulkloadFile
BulkloadFile
类,初始化 Lakehouse 连接和 BulkloadStreamV2 对象。
- 读取本地 CSV 文件并写入 Lakehouse:使用 Java IO 流读取本地 CSV 文件,并将数据逐行写入 Lakehouse。
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.io.BufferedReader;
import java.io.File;
import java.io.FileReader;
import java.text.MessageFormat;
import java.util.Arrays;
import java.util.Collection;
public class BulkloadFile {
private static ClickZettaClient client;
private static final String password = "";
private static final String table = "bulk_order_items";
private static final String workspace = "";
private static final String schema = "public";
private static final String vc = "default";
private static final String user = "";
static BulkloadStreamV2 bulkloadStream;
public static void main(String[] args) throws Exception {
try {
initialize();
File csvFile = new File("olist_order_items_dataset.csv");
try (BufferedReader reader = new BufferedReader(new FileReader(csvFile))) {
// Skip the header row
reader.readLine(); // Skip the first line (header)
// Insert data into the database
String line;
while ((line = reader.readLine()) != null) {
String[] values = line.split(",");
// 类型转化保持和服务端类型一致
String orderId = values[0];
int orderItemId = Integer.parseInt(values[1]); // Convert order_item_id to int
String productId = values[2];
String sellerId = values[3];
String shippingLimitDate = values[4];
double price = Double.parseDouble(values[5]);
double freightValue = Double.parseDouble(values[6]);
Row row = bulkloadStream.createRow(0);
// Set parameter values
row.setValue(0, orderId);
row.setValue(1, orderItemId);
row.setValue(2, productId);
row.setValue(3, sellerId);
row.setValue(4, shippingLimitDate);
row.setValue(5, price);
row.setValue(6, freightValue);
// 必须调用该方法,否则无法将数据发送到服务端
bulkloadStream.apply(row, 0);
}
}
// 提交本次写入:收集提交请求 -> 预提交 -> 提交 -> 等待完成。
Collection<BulkLoadCommitter.CommitRequest<Committable>> commitRequests =
bulkloadStream.getCommitRequests();
String transactionId = bulkloadStream.prepareCommit(commitRequests);
bulkloadStream.commit(Arrays.asList(transactionId), commitRequests).get();
System.out.println("Data inserted successfully!");
} finally {
if (bulkloadStream != null) {
bulkloadStream.close();
}
if (client != null) {
client.close();
}
}
}
private static void initialize() throws Exception {
String url = MessageFormat.format("jdbc:clickzetta://demo_instance.cn-shanghai-alicloud.api.clickzetta.com/{0}?" +
"schema={1}&username={2}&password={3}&virtualcluster={4}",
workspace, schema, user, password, vc);
client = ClickZettaClient.newBuilder().url(url).build();
bulkloadStream = client.newBulkloadStreamV2Builder()
.withOperation(BulkLoadOperation.APPEND)
.createStream(client, schema, table);
}
}