Use sync insert strategy since we are already batching client-side (#2303)

This commit is contained in:
Eric Allam
2025-07-23 17:17:13 +01:00
committed by GitHub
parent ed86b4f5db
commit b447a8040f
4 changed files with 21 additions and 10 deletions
+1
View File
@@ -911,6 +911,7 @@ const EnvironmentSchema = z.object({
RUN_REPLICATION_INSERT_MAX_RETRIES: z.coerce.number().int().default(3),
RUN_REPLICATION_INSERT_BASE_DELAY_MS: z.coerce.number().int().default(100),
RUN_REPLICATION_INSERT_MAX_DELAY_MS: z.coerce.number().int().default(2000),
RUN_REPLICATION_INSERT_STRATEGY: z.enum(["insert", "insert_async"]).default("insert"),
// Clickhouse
CLICKHOUSE_URL: z.string(),
@@ -65,6 +65,7 @@ function initializeRunsReplicationInstance() {
insertMaxRetries: env.RUN_REPLICATION_INSERT_MAX_RETRIES,
insertBaseDelayMs: env.RUN_REPLICATION_INSERT_BASE_DELAY_MS,
insertMaxDelayMs: env.RUN_REPLICATION_INSERT_MAX_DELAY_MS,
insertStrategy: env.RUN_REPLICATION_INSERT_STRATEGY,
});
if (env.RUN_REPLICATION_ENABLED === "1") {
@@ -53,6 +53,7 @@ export type RunsReplicationServiceOptions = {
logLevel?: LogLevel;
tracer?: Tracer;
waitForAsyncInsert?: boolean;
insertStrategy?: "insert" | "insert_async";
// Retry configuration for insert operations
insertMaxRetries?: number;
insertBaseDelayMs?: number;
@@ -90,6 +91,7 @@ export class RunsReplicationService {
private _insertMaxRetries: number;
private _insertBaseDelayMs: number;
private _insertMaxDelayMs: number;
private _insertStrategy: "insert" | "insert_async";
public readonly events: EventEmitter<RunsReplicationServiceEvents>;
@@ -101,6 +103,8 @@ export class RunsReplicationService {
this._acknowledgeTimeoutMs = options.acknowledgeTimeoutMs ?? 1_000;
this._insertStrategy = options.insertStrategy ?? "insert";
this._replicationClient = new LogicalReplicationClient({
pgConfig: {
connectionString: options.pgConnectionUrl,
@@ -598,15 +602,26 @@ export class RunsReplicationService {
return delay + jitter;
}
#getClickhouseInsertSettings() {
if (this._insertStrategy === "insert") {
return {};
} else if (this._insertStrategy === "insert_async") {
return {
async_insert: 1 as const,
async_insert_max_data_size: "1000000",
async_insert_busy_timeout_ms: 1000,
wait_for_async_insert: this.options.waitForAsyncInsert ? (1 as const) : (0 as const),
};
}
}
async #insertTaskRunInserts(taskRunInserts: TaskRunV2[], attempt: number) {
return await startSpan(this._tracer, "insertTaskRunsInserts", async (span) => {
const [insertError, insertResult] = await this.options.clickhouse.taskRuns.insert(
taskRunInserts,
{
params: {
clickhouse_settings: {
wait_for_async_insert: this.options.waitForAsyncInsert ? 1 : 0,
},
clickhouse_settings: this.#getClickhouseInsertSettings(),
},
}
);
@@ -631,9 +646,7 @@ export class RunsReplicationService {
payloadInserts,
{
params: {
clickhouse_settings: {
wait_for_async_insert: this.options.waitForAsyncInsert ? 1 : 0,
},
clickhouse_settings: this.#getClickhouseInsertSettings(),
},
}
);
@@ -56,10 +56,6 @@ export function insertTaskRuns(ch: ClickhouseWriter, settings?: ClickHouseSettin
table: "trigger_dev.task_runs_v2",
schema: TaskRunV2,
settings: {
async_insert: 1,
wait_for_async_insert: 0,
async_insert_max_data_size: "1000000",
async_insert_busy_timeout_ms: 1000,
enable_json_type: 1,
type_json_skip_duplicated_paths: 1,
...settings,