Пакетная вставка данных
YDB поддерживает пакетную вставку большого количества строк без гарантий атомарности. Запись данных разбивается на несколько независимых транзакций, каждая их которых затрагивает единственную партицию, с параллельным исполнением. За счет этого такой подход более эффективен чем YQL. В случае успеха метод BulkUpsert гарантирует вставку всех данных, переданных в рамках данного запроса.
Важно
При использовании BulkUpsert для вставки данных в колоночные таблицы необходимо передавать значения всех колонок, включая NULL-значения.
Ниже приведены примеры кода использования встроенных в YDB SDK средств выполнения пакетной вставки:
#include <ydb-cpp-sdk/client/table/table.h>
void BulkUpsertLogs(const NYdb::TDriver& driver) {
NYdb::NTable::TTableClient client(driver);
constexpr int kBatchSize = 1000;
NYdb::TValueBuilder rowsBuilder;
rowsBuilder.BeginList();
for (int i = 0; i < kBatchSize; ++i) {
rowsBuilder.AddListItem()
.BeginStruct()
.AddMember("App").Utf8("App_" + std::to_string(i / 256))
.AddMember("Host").Utf8("192.168.0." + std::to_string(i % 256))
.AddMember("Timestamp")
.Timestamp(TInstant::Now() + TDuration::Seconds(i))
.AddMember("HttpCode").Uint32(static_cast<uint32_t>(i % 113 == 0 ? 404 : 200))
.AddMember("Message")
.Utf8(i % 3 == 0 ? "GET / HTTP/1.1" : "GET /images/logo.png HTTP/1.1")
.EndStruct();
}
rowsBuilder.EndList();
NYdb::TValue rows = rowsBuilder.Build();
NYdb::NStatusHelpers::ThrowOnError(client.RetryOperationSync(
[&rows](NYdb::NTable::TTableClient& client) {
return client.BulkUpsert("/local/bulk_upsert_example", NYdb::TValue{rows}).GetValueSync();
},
NYdb::NTable::TRetryOperationSettings()
.Idempotent(true)
));
}
#include <userver/ydb/io/supported_types.hpp>
#include <userver/ydb/table.hpp>
struct LogMessage final {
ydb::Utf8 App;
ydb::Utf8 Host;
std::chrono::system_clock::time_point Timestamp;
std::uint32_t HttpCode;
ydb::Utf8 Message;
};
void BulkUpsertLogs(ydb::TableClient& client) {
constexpr int kBatchSize = 1000;
std::vector<LogMessage> rows;
rows.reserve(kBatchSize);
for (int i = 0; i < kBatchSize; ++i) {
rows.push_back({
.App = ydb::Utf8{"App_" + std::to_string(i / 256)},
.Host = ydb::Utf8{"192.168.0." + std::to_string(i % 256)},
.Timestamp = std::chrono::system_clock::now() + std::chrono::seconds{i},
.HttpCode = static_cast<std::uint32_t>(i % 113 == 0 ? 404 : 200),
.Message = ydb::Utf8{
i % 3 == 0 ? "GET / HTTP/1.1" : "GET /images/logo.png HTTP/1.1"
},
});
}
client.BulkUpsert(
"/local/bulk_upsert_example",
rows,
ydb::OperationSettings{
.is_idempotent = true,
}
);
}
Пакетная вставка нативных YDB данных
package main
import (
"context"
"os"
"github.com/ydb-platform/ydb-go-sdk/v3"
"github.com/ydb-platform/ydb-go-sdk/v3/table"
"github.com/ydb-platform/ydb-go-sdk/v3/table/types"
)
func main() {
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
db, err := ydb.Open(ctx,
os.Getenv("YDB_CONNECTION_STRING"),
ydb.WithAccessTokenCredentials(os.Getenv("YDB_TOKEN")),
)
if err != nil {
panic(err)
}
defer db.Close(ctx)
type logMessage struct {
App string
Host string
Timestamp time.Time
HTTPCode uint32
Message string
}
// prepare native go data
const batchSize = 10000
logs := make([]logMessage, 0, batchSize)
for i := 0; i < batchSize; i++ {
message := logMessage{
App: fmt.Sprintf("App_%d", i/256),
Host: fmt.Sprintf("192.168.0.%d", i%256),
Timestamp: time.Now().Add(time.Millisecond * time.Duration(i%1000)),
HTTPCode: 200,
}
if i%2 == 0 {
message.Message = "GET / HTTP/1.1"
} else {
message.Message = "GET /images/logo.png HTTP/1.1"
}
logs = append(logs, message)
}
// execute bulk upsert with native ydb data
err = db.Table().Do( // Do retry operation on errors with best effort
ctx, // context manage exiting from Do
func(ctx context.Context, s table.Session) (err error) { // retry operation
rows := make([]types.Value, 0, len(logs))
for _, msg := range logs {
rows = append(rows, types.StructValue(
types.StructFieldValue("App", types.UTF8Value(msg.App)),
types.StructFieldValue("Host", types.UTF8Value(msg.Host)),
types.StructFieldValue("Timestamp", types.TimestampValueFromTime(msg.Timestamp)),
types.StructFieldValue("HTTPCode", types.Uint32Value(msg.HTTPCode)),
types.StructFieldValue("Message", types.UTF8Value(msg.Message)),
))
}
return s.BulkUpsert(ctx, "/local/bulk_upsert_example", types.ListValue(rows...))
},
)
if err != nil {
fmt.Printf("unexpected error: %v", err)
}
}
Пакетная вставка CSV
package main
import (
"context"
"fmt"
"os"
"github.com/ydb-platform/ydb-go-sdk/v3"
"github.com/ydb-platform/ydb-go-sdk/v3/table"
)
func main() {
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
db, err := ydb.Open(ctx,
os.Getenv("YDB_CONNECTION_STRING"),
ydb.WithAccessTokenCredentials(os.Getenv("YDB_TOKEN")),
)
if err != nil {
panic(err)
}
defer db.Close(ctx)
csv := `skip row
id,val
42,"text42"
43,"text43"
44,hello
`
err = db.Table().BulkUpsert(ctx, "/local/bulk_upsert_example", table.BulkUpsertDataCsv(
[]byte(csv),
table.WithCsvHeader(),
table.WithCsvSkipRows(2),
table.WithCsvNullValue([]byte("hello")), // строка "hello" будет восприниматься как NULL
))
if err != nil {
fmt.Printf("unexpected error: %v", err)
}
}
Пакетная вставка Apache Arrow
В следующем примере для подготовки данных используется пакет arrow.
package main
import (
"bytes"
"context"
"fmt"
"github.com/apache/arrow-go/v18/arrow"
"github.com/apache/arrow-go/v18/arrow/array"
"github.com/apache/arrow-go/v18/arrow/ipc"
"github.com/apache/arrow-go/v18/arrow/memory"
"github.com/ydb-platform/ydb-go-sdk/v3"
"github.com/ydb-platform/ydb-go-sdk/v3/table"
)
func main() {
ctx := context.Background()
db, err := ydb.Open(ctx,
os.Getenv("YDB_CONNECTION_STRING"),
ydb.WithAccessTokenCredentials(os.Getenv("YDB_TOKEN")),
)
if err != nil {
panic(err)
}
defer db.Close(ctx) // cleanup resources
mem := memory.NewGoAllocator()
schema := arrow.NewSchema([]arrow.Field{
{Name: "id", Type: arrow.PrimitiveTypes.Int64},
{Name: "val", Type: arrow.BinaryTypes.String},
}, nil)
b := array.NewRecordBuilder(mem, schema)
defer b.Release()
b.Field(0).(*array.Int64Builder).AppendValues(
[]int64{123, 234}, nil)
b.Field(1).(*array.StringBuilder).AppendValues(
[]string{"data1", "data2"}, nil)
rec := b.NewRecordBatch()
defer rec.Release()
schemaPayload := ipc.GetSchemaPayload(rec.Schema(), mem)
defer schemaPayload.Release()
dataPayload, err := ipc.GetRecordBatchPayload(rec)
if err != nil {
panic(err)
}
defer dataPayload.Release()
var schemaBuf bytes.Buffer
_, err = schemaPayload.WritePayload(&schemaBuf)
if err != nil {
panic(err)
}
var dataBuf bytes.Buffer
_, err = dataPayload.WritePayload(&dataBuf)
if err != nil {
panic(err)
}
err = db.Table().BulkUpsert(ctx, "/local/bulk_upsert_example", table.BulkUpsertDataArrow(
dataBuf.Bytes(),
table.WithArrowSchema(schemaBuf.Bytes()), // schema is required
))
if err != nil {
fmt.Printf("unexpected error: %v", err)
}
}
Реализация database/sql драйвера для YDB не поддерживает нетранзакционную пакетную вставку данных.
Для пакетной вставки следует пользоваться транзакционной вставкой.
Пакетная вставка через BulkUpsert эффективнее транзакционного YQL для больших объёмов данных. Для небольших наборов строк см. UPSERT. Структура таблицы описана в разделе Таблицы.
import java.time.Instant;
import java.util.ArrayList;
import java.util.List;
import tech.ydb.auth.NopAuthProvider;
import tech.ydb.common.transaction.TxMode;
import tech.ydb.core.grpc.GrpcTransport;
import tech.ydb.query.QueryClient;
import tech.ydb.query.result.ResultSetReader;
import tech.ydb.query.tools.QueryReader;
import tech.ydb.query.tools.SessionRetryContext;
import tech.ydb.table.TableClient;
import tech.ydb.table.query.Params;
import tech.ydb.table.settings.BulkUpsertSettings;
import tech.ydb.table.values.ListType;
import tech.ydb.table.values.ListValue;
import tech.ydb.table.values.PrimitiveType;
import tech.ydb.table.values.PrimitiveValue;
import tech.ydb.table.values.StructType;
import tech.ydb.table.values.Value;
public class BulkUpsertExample {
private static final String TABLE_NAME = "bulk_upsert";
private static final int BATCH_SIZE = 1000;
public static void main(String[] args) {
String connectionString = System.getenv().getOrDefault(
"YDB_CONNECTION_STRING", "grpc://localhost:2136/local");
try (GrpcTransport transport = GrpcTransport.forConnectionString(connectionString)
.withAuthProvider(NopAuthProvider.INSTANCE)
.build();
QueryClient queryClient = QueryClient.newClient(transport).build();
TableClient tableClient = TableClient.newClient(transport).build()) {
SessionRetryContext queryRetry = SessionRetryContext.create(queryClient).build();
SessionRetryContext tableRetry = SessionRetryContext.create(tableClient).build();
// Создаём таблицу для пакетной вставки
queryRetry.supplyResult(session -> QueryReader.readFrom(session.createQuery("""
CREATE TABLE IF NOT EXISTS bulk_upsert (
app Text,
timestamp Timestamp,
host Text,
http_code Uint32,
message Text,
PRIMARY KEY (app, timestamp, host)
);
""", TxMode.NONE, Params.empty())
)).join().getValue();
// Полный путь к таблице: /local/bulk_upsert
String tablePath = transport.getDatabase() + "/" + TABLE_NAME;
StructType rowType = StructType.of(
"app", PrimitiveType.Text,
"timestamp", PrimitiveType.Timestamp,
"host", PrimitiveType.Text,
"http_code", PrimitiveType.Uint32,
"message", PrimitiveType.Text
);
// Генерация пакета записей
List<Value<?>> rows = new ArrayList<>(BATCH_SIZE);
for (int i = 0; i < BATCH_SIZE; i++) {
rows.add(rowType.newValue(
"app", PrimitiveValue.newText("App_" + i / 256),
"timestamp", PrimitiveValue.newTimestamp(Instant.now().plusSeconds(i)),
"host", PrimitiveValue.newText("192.168.0." + i % 256),
"http_code", PrimitiveValue.newUint32(i % 113 == 0 ? 404 : 200),
"message", PrimitiveValue.newText(
i % 3 == 0 ? "GET / HTTP/1.1" : "GET /images/logo.png HTTP/1.1")
));
}
ListValue batch = ListType.of(rowType).newValue(rows);
// Пакетная вставка без атомарности всего батча
tableRetry.supplyStatus(session ->
session.executeBulkUpsert(tablePath, batch, new BulkUpsertSettings())
).join().expectSuccess("bulk upsert failed");
// Проверка количества строк
QueryReader countReader = queryRetry.supplyResult(session -> QueryReader.readFrom(
session.createQuery(
"SELECT COUNT(*) AS cnt FROM bulk_upsert", TxMode.NONE, Params.empty())
)).join().getValue();
ResultSetReader rs = countReader.getResultSet(0);
if (rs.next()) {
System.out.println("Строк в таблице bulk_upsert: " + rs.getColumn("cnt").getUint64());
}
}
}
}
import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.sql.Statement;
import java.sql.Timestamp;
import java.time.Instant;
public class JdbcBulkUpsertExample {
private static final int BATCH_SIZE = 1000;
public static void main(String[] args) throws SQLException {
String url = System.getenv().getOrDefault(
"YDB_JDBC_URL", "jdbc:ydb:grpc://localhost:2136/local");
try (Connection conn = DriverManager.getConnection(url)) {
try (Statement stmt = conn.createStatement()) {
stmt.execute("""
CREATE TABLE IF NOT EXISTS bulk_upsert (
app Text,
timestamp Timestamp,
host Text,
http_code Uint32,
message Text,
PRIMARY KEY (app, timestamp, host)
);
""");
}
String bulkSql = """
BULK UPSERT INTO bulk_upsert (app, timestamp, host, http_code, message)
VALUES (?, ?, ?, ?, ?);
""";
try (PreparedStatement ps = conn.prepareStatement(bulkSql)) {
for (int i = 0; i < BATCH_SIZE; i++) {
ps.setString(1, "App_" + i / 256);
ps.setTimestamp(2, Timestamp.from(Instant.now().plusSeconds(i)));
ps.setString(3, "192.168.0." + i % 256);
ps.setLong(4, i % 113 == 0 ? 404 : 200);
ps.setString(5, i % 3 == 0 ? "GET / HTTP/1.1" : "GET /images/logo.png HTTP/1.1");
ps.addBatch();
}
ps.executeBatch();
}
try (Statement stmt = conn.createStatement();
ResultSet rs = stmt.executeQuery("SELECT COUNT(*) AS cnt FROM bulk_upsert")) {
if (rs.next()) {
System.out.println("Строк в таблице bulk_upsert: " + rs.getLong("cnt"));
}
}
}
}
}
В Spring Boot, Hibernate, JOOQ и других фреймворках вокруг ORM поверх JDBC можно выполнять нативный YQL (в том числе из репозиториев и @Query). Драйвер стремится оптимизировать крупные вставки; операции UPDATE, INSERT, DELETE, UPSERT, идущие через JDBC, при необходимости автоматически группируются в пакеты на стороне драйвера.
import posixpath
import ydb
def bulk_upsert(driver: ydb.Driver, path: str):
column_types = (
ydb.BulkUpsertColumns()
.add_column("id", ydb.PrimitiveType.Uint64)
.add_column("val", ydb.OptionalType(ydb.PrimitiveType.Utf8))
)
rows = [
{"id": 1, "val": "1"},
{"id": 2, "val": "2"},
{"id": 3, "val": "3"},
]
driver.table_client.bulk_upsert(posixpath.join(path, "tablename"), rows, column_types)
import os
import posixpath
import ydb
import asyncio
async def bulk_upsert(driver: ydb.aio.Driver, path: str):
column_types = (
ydb.BulkUpsertColumns()
.add_column("id", ydb.PrimitiveType.Uint64)
.add_column("val", ydb.OptionalType(ydb.PrimitiveType.Utf8))
)
rows = [
{"id": 1, "val": "1"},
{"id": 2, "val": "2"},
{"id": 3, "val": "3"},
]
await driver.table_client.bulk_upsert(
posixpath.join(path, "tablename"), rows, column_types
)
async def main():
async with ydb.aio.Driver(
connection_string=os.environ["YDB_CONNECTION_STRING"],
credentials=ydb.credentials_from_env_variables(),
) as driver:
await driver.wait()
await bulk_upsert(driver, "/local")
asyncio.run(main())
import os
import sqlalchemy as sa
import ydb
engine = sa.create_engine(os.environ["YDB_SQLALCHEMY_URL"])
with engine.connect() as connection:
dbapi_conn = connection.connection
column_types = (
ydb.BulkUpsertColumns()
.add_column("id", ydb.PrimitiveType.Uint64)
.add_column("val", ydb.OptionalType(ydb.PrimitiveType.Utf8))
)
rows = [
{"id": 1, "val": "1"},
{"id": 2, "val": "2"},
{"id": 3, "val": "3"},
]
dbapi_conn.bulk_upsert("tablename", rows, column_types)
Функциональность на данный момент не поддерживается.
Функциональность на данный момент не поддерживается.
use ydb::{ydb_struct, AccessTokenCredentials, ClientBuilder, Value, YdbResult};
#[tokio::main]
async fn main() -> YdbResult<()> {
let client = ClientBuilder::new_from_connection_string(
"grpc://localhost:2136?database=local",
)?
.with_credentials(AccessTokenCredentials::from("..."))
.client()?;
client.wait().await?;
let rows: Vec<Value> = vec![
ydb_struct!(
"id" => 1_u64,
"val" => Value::Text("1".into()),
),
ydb_struct!(
"id" => 2_u64,
"val" => Value::Text("2".into()),
),
ydb_struct!(
"id" => 3_u64,
"val" => Value::Text("3".into()),
),
];
client
.table_client()
.bulk_upsert("/local/tablename", rows)
.idempotent(true)
.await?;
Ok(())
}
<?php
use YdbPlatform\Ydb\Ydb;
$ydb = new Ydb($config);
$rows = [
['id' => 1, 'val' => '1'],
['id' => 2, 'val' => '2'],
['id' => 3, 'val' => '3'],
];
$ydb->table()->bulkUpsert('tablename', $rows, [
'id' => 'UINT64',
'val' => 'UTF8',
]);