feat(etl): Add new setting for controlling table copy (#42853)

This commit is contained in:
Riccardo Busetti authored and GitHub committed 2026-02-24 08:14:23 +01:00
1 parent 2b4036ff7b
commit 67ed8f18ee
8 files changed
+145 -88

No files matched your search

@@ -72,35 +72,6 @@ export const AdvancedSettings = ({
)}
/>
<FormField_Shadcn_
control={form.control}
name="maxSize"
render={({ field }) => (
<FormItemLayout
label="Batch size"
layout="horizontal"
description={
<>
<p>Number of rows to send in a batch.</p>
<p>Larger batches use more memory, with the risk of running out of memory.</p>
</>
}
>
<FormControl_Shadcn_>
<PrePostTab postTab="rows">
<Input_Shadcn_
{...field}
type="number"
value={field.value ?? ''}
onChange={handleNumberChange(field)}
placeholder="Default: 100000"
/>
</PrePostTab>
</FormControl_Shadcn_>
</FormItemLayout>
)}
/>
<FormField_Shadcn_
control={form.control}
name="maxTableSyncWorkers"
@@ -130,6 +101,37 @@ export const AdvancedSettings = ({
)}
/>
<FormField_Shadcn_
control={form.control}
name="maxCopyConnectionsPerTable"
render={({ field }) => (
<FormItemLayout
label="Copy connections per table"
layout="horizontal"
description={
<>
<p>
Number of copy connections each table sync worker can use in parallel when
copying a table.
</p>
</>
}
>
<FormControl_Shadcn_>
<PrePostTab postTab="connections">
<Input_Shadcn_
{...field}
type="number"
value={field.value ?? ''}
onChange={handleNumberChange(field)}
placeholder="Default: 4"
/>
</PrePostTab>
</FormControl_Shadcn_>
</FormItemLayout>
)}
/>
{/* BigQuery-specific: Max staleness */}
{type === 'BigQuery' && (
<FormField_Shadcn_
@@ -7,12 +7,16 @@ export const DestinationPanelFormSchema = z.object({
name: z.string().min(1, 'Name is required'),
publicationName: z.string().min(1, 'Publication is required'),
maxFillMs: z.number().min(1, 'Max Fill milliseconds should be greater than 0').int().optional(),
maxSize: z.number().min(1, 'Max batch size should be greater than 0').int().optional(),
maxTableSyncWorkers: z
.number()
.min(1, 'Max table sync workers should be greater than 0')
.int()
.optional(),
maxCopyConnectionsPerTable: z
.number()
.int()
.min(1, 'Max copy connections per table should be greater than 0')
.optional(),
// BigQuery fields
projectId: z.string().optional(),
datasetId: z.string().optional(),
@@ -19,8 +19,8 @@ import { useReplicationSourcesQuery } from 'data/replication/sources-query'
import { useStartPipelineMutation } from 'data/replication/start-pipeline-mutation'
import { useUpdateDestinationPipelineMutation } from 'data/replication/update-destination-pipeline-mutation'
import {
type ValidationFailure,
useValidateDestinationMutation,
type ValidationFailure,
} from 'data/replication/validate-destination-mutation'
import { useValidatePipelineMutation } from 'data/replication/validate-pipeline-mutation'
import { useIcebergNamespaceCreateMutation } from 'data/storage/iceberg-namespace-create-mutation'
@@ -174,8 +174,8 @@ export const DestinationForm = ({
name: destinationData?.name ?? '',
publicationName: pipelineData?.config.publication_name ?? '',
maxFillMs: pipelineData?.config?.batch?.max_fill_ms ?? undefined,
maxSize: pipelineData?.config?.batch?.max_size ?? undefined,
maxTableSyncWorkers: pipelineData?.config?.max_table_sync_workers ?? undefined,
maxCopyConnectionsPerTable: pipelineData?.config?.max_copy_connections_per_table ?? undefined,
// BigQuery fields
projectId: isBigQueryConfig ? config.big_query.project_id : '',
datasetId: isBigQueryConfig ? config.big_query.dataset_id : '',
@@ -303,8 +303,8 @@ export const DestinationForm = ({
sourceId,
publicationName: data.publicationName,
maxFillMs: data.maxFillMs,
maxSize: data.maxSize,
maxTableSyncWorkers: data.maxTableSyncWorkers,
maxCopyConnectionsPerTable: data.maxCopyConnectionsPerTable,
}),
])
@@ -403,10 +403,9 @@ export const DestinationForm = ({
}
const batchConfig: BatchConfig | undefined =
data.maxFillMs !== undefined || data.maxSize !== undefined
data.maxFillMs !== undefined
? {
...(data.maxFillMs !== undefined ? { maxFillMs: data.maxFillMs } : {}),
...(data.maxSize !== undefined ? { maxSize: data.maxSize } : {}),
}
: undefined
const hasBatchFields = batchConfig !== undefined
@@ -422,6 +421,7 @@ export const DestinationForm = ({
pipelineConfig: {
publicationName: data.publicationName,
maxTableSyncWorkers: data.maxTableSyncWorkers,
maxCopyConnectionsPerTable: data.maxCopyConnectionsPerTable,
...(hasBatchFields ? { batch: batchConfig } : {}),
},
sourceId,
@@ -486,10 +486,9 @@ export const DestinationForm = ({
destinationConfig = { iceberg: icebergConfig }
}
const batchConfig: BatchConfig | undefined =
data.maxFillMs !== undefined || data.maxSize !== undefined
data.maxFillMs !== undefined
? {
...(data.maxFillMs !== undefined ? { maxFillMs: data.maxFillMs } : {}),
...(data.maxSize !== undefined ? { maxSize: data.maxSize } : {}),
}
: undefined
const hasBatchFields = batchConfig !== undefined
@@ -504,6 +503,7 @@ export const DestinationForm = ({
pipelineConfig: {
publicationName: data.publicationName,
maxTableSyncWorkers: data.maxTableSyncWorkers,
maxCopyConnectionsPerTable: data.maxCopyConnectionsPerTable,
...(hasBatchFields ? { batch: batchConfig } : {}),
},
})
@@ -1,9 +1,9 @@
import { useMutation, useQueryClient } from '@tanstack/react-query'
import { toast } from 'sonner'
import type { components } from 'api-types'
import { handleError, post } from 'data/fetchers'
import { toast } from 'sonner'
import type { ResponseError, UseCustomMutationOptions } from 'types'
import { replicationKeys } from './keys'
export type DestinationConfig =
@@ -33,7 +33,6 @@ export type IcebergDestinationConfig = {
export type BatchConfig = {
maxFillMs?: number
maxSize?: number
}
export type CreateDestinationPipelineParams = {
@@ -45,6 +44,7 @@ export type CreateDestinationPipelineParams = {
publicationName: string
batch?: BatchConfig
maxTableSyncWorkers?: number
maxCopyConnectionsPerTable?: number
}
}
@@ -53,7 +53,7 @@ async function createDestinationPipeline(
projectRef,
destinationName: destinationName,
destinationConfig,
pipelineConfig: { publicationName, batch, maxTableSyncWorkers },
pipelineConfig: { publicationName, batch, maxTableSyncWorkers, maxCopyConnectionsPerTable },
sourceId,
}: CreateDestinationPipelineParams,
signal?: AbortSignal
@@ -113,11 +113,13 @@ async function createDestinationPipeline(
...(maxTableSyncWorkers !== undefined
? { max_table_sync_workers: maxTableSyncWorkers }
: {}),
...(maxCopyConnectionsPerTable !== undefined
? { max_copy_connections_per_table: maxCopyConnectionsPerTable }
: {}),
...(batch
? {
batch: {
...(batch.maxFillMs !== undefined ? { max_fill_ms: batch.maxFillMs } : {}),
...(batch.maxSize !== undefined ? { max_size: batch.maxSize } : {}),
},
}
: {}),
@@ -1,9 +1,9 @@
import { useMutation, useQueryClient } from '@tanstack/react-query'
import { toast } from 'sonner'
import type { components } from 'api-types'
import { handleError, post } from 'data/fetchers'
import { toast } from 'sonner'
import type { ResponseError, UseCustomMutationOptions } from 'types'
import { BatchConfig, DestinationConfig } from './create-destination-pipeline-mutation'
import { replicationKeys } from './keys'
@@ -18,6 +18,7 @@ export type UpdateDestinationPipelineParams = {
publicationName: string
batch?: BatchConfig
maxTableSyncWorkers?: number
maxCopyConnectionsPerTable?: number
}
}
@@ -28,7 +29,7 @@ async function updateDestinationPipeline(
projectRef,
destinationName: destinationName,
destinationConfig,
pipelineConfig: { publicationName, batch, maxTableSyncWorkers },
pipelineConfig: { publicationName, batch, maxTableSyncWorkers, maxCopyConnectionsPerTable },
sourceId,
}: UpdateDestinationPipelineParams,
signal?: AbortSignal
@@ -88,11 +89,13 @@ async function updateDestinationPipeline(
...(maxTableSyncWorkers !== undefined
? { max_table_sync_workers: maxTableSyncWorkers }
: {}),
...(maxCopyConnectionsPerTable !== undefined
? { max_copy_connections_per_table: maxCopyConnectionsPerTable }
: {}),
...(batch
? {
batch: {
...(batch.maxFillMs !== undefined ? { max_fill_ms: batch.maxFillMs } : {}),
...(batch.maxSize !== undefined ? { max_size: batch.maxSize } : {}),
},
}
: {}),
@@ -9,8 +9,8 @@ type ValidatePipelineParams = {
sourceId: number
publicationName: string
maxFillMs?: number
maxSize?: number
maxTableSyncWorkers?: number
maxCopyConnectionsPerTable?: number
}
type ValidatePipelineResponse = components['schemas']['ValidatePipelineResponse']
@@ -20,18 +20,15 @@ async function validatePipeline(
sourceId,
publicationName,
maxFillMs,
maxSize,
maxTableSyncWorkers,
maxCopyConnectionsPerTable,
}: ValidatePipelineParams,
signal?: AbortSignal
): Promise<ValidatePipelineResponse> {
if (!projectRef) throw new Error('projectRef is required')
if (!sourceId) throw new Error('sourceId is required')
const batchConfig =
maxFillMs !== undefined || maxSize !== undefined
? { max_fill_ms: maxFillMs, max_size: maxSize }
: undefined
const batchConfig = maxFillMs !== undefined ? { max_fill_ms: maxFillMs } : undefined
const { data, error } = await post('/platform/replication/{ref}/pipelines/validate', {
params: { path: { ref: projectRef } },
@@ -40,6 +37,7 @@ async function validatePipeline(
config: {
publication_name: publicationName,
max_table_sync_workers: maxTableSyncWorkers,
max_copy_connections_per_table: maxCopyConnectionsPerTable,
batch: batchConfig,
},
},
+70 -1
View File
@@ -934,7 +934,8 @@ export interface paths {
path?: never
cookie?: never
}
get?: never
/** Get database disk attributes */
get: operations['v1-get-database-disk']
put?: never
/** Modify database disk */
post: operations['v1-modify-database-disk']
@@ -2806,6 +2807,23 @@ export interface components {
type: 'io2'
}
}
DiskResponse: {
attributes:
| {
iops: number
size_gb: number
throughput_mibps?: number
/** @enum {string} */
type: 'gp3'
}
| {
iops: number
size_gb: number
/** @enum {string} */
type: 'io2'
}
last_modified_at?: string
}
DiskUtilMetricsResponse: {
metrics: {
fs_avail_bytes: number
@@ -4821,6 +4839,7 @@ export interface components {
* @description Deprecated. Use `status` instead.
*/
healthy: boolean
replication_connected: boolean
}
| {
db_schema: string
@@ -8291,6 +8310,56 @@ export interface operations {
}
}
}
'v1-get-database-disk': {
parameters: {
query?: never
header?: never
path: {
/** @description Project ref */
ref: string
}
cookie?: never
}
requestBody?: never
responses: {
200: {
headers: {
[name: string]: unknown
}
content: {
'application/json': components['schemas']['DiskResponse']
}
}
/** @description Unauthorized */
401: {
headers: {
[name: string]: unknown
}
content?: never
}
/** @description Forbidden action */
403: {
headers: {
[name: string]: unknown
}
content?: never
}
/** @description Rate limit exceeded */
429: {
headers: {
[name: string]: unknown
}
content?: never
}
/** @description Failed to get database disk attributes */
500: {
headers: {
[name: string]: unknown
}
content?: never
}
}
}
'v1-modify-database-disk': {
parameters: {
query?: never
+14 -35
View File
@@ -5451,12 +5451,9 @@ export interface components {
* @example 200
*/
max_fill_ms?: number
/**
* @description Maximum size of the batch
* @example 200
*/
max_size?: number
}
/** @description Maximum number of copy connections per table */
max_copy_connections_per_table?: number
/** @description Maximum number of table sync workers */
max_table_sync_workers?: number
/**
@@ -5481,12 +5478,9 @@ export interface components {
* @example 200
*/
max_fill_ms?: number
/**
* @description Maximum size of the batch
* @example 200
*/
max_size?: number
}
/** @description Maximum number of copy connections per table */
max_copy_connections_per_table?: number
/** @description Maximum number of table sync workers */
max_table_sync_workers?: number
/**
@@ -8766,12 +8760,9 @@ export interface components {
* @example 200
*/
max_fill_ms?: number
/**
* @description Maximum size of the batch
* @example 200
*/
max_size?: number
}
/** @description Maximum number of copy connections per table */
max_copy_connections_per_table?: number
/** @description Maximum number of table sync workers */
max_table_sync_workers?: number
/**
@@ -8828,12 +8819,9 @@ export interface components {
* @example 200
*/
max_fill_ms?: number
/**
* @description Maximum size of the batch
* @example 200
*/
max_size?: number
}
/** @description Maximum number of copy connections per table */
max_copy_connections_per_table?: number
/** @description Maximum number of table sync workers */
max_table_sync_workers?: number
/**
@@ -10365,12 +10353,9 @@ export interface components {
* @example 200
*/
max_fill_ms?: number
/**
* @description Maximum size of the batch
* @example 200
*/
max_size?: number
}
/** @description Maximum number of copy connections per table */
max_copy_connections_per_table?: number
/** @description Maximum number of table sync workers */
max_table_sync_workers?: number
/**
@@ -10395,12 +10380,9 @@ export interface components {
* @example 200
*/
max_fill_ms?: number
/**
* @description Maximum size of the batch
* @example 200
*/
max_size?: number
}
/** @description Maximum number of copy connections per table */
max_copy_connections_per_table?: number
/** @description Maximum number of table sync workers */
max_table_sync_workers?: number
/**
@@ -10806,12 +10788,9 @@ export interface components {
* @example 200
*/
max_fill_ms?: number
/**
* @description Maximum size of the batch
* @example 200
*/
max_size?: number
}
/** @description Maximum number of copy connections per table */
max_copy_connections_per_table?: number
/** @description Maximum number of table sync workers */
max_table_sync_workers?: number
/**