Skip to content
Query.Farm
Talk with Us

Buffering functions

On this page

Sink, combine, source — for output that depends on the whole input.

source
export function defineTableBufferingFunction<
TArgs = Record<string, any>,
TState = any,
>(config: TableBufferingConfig<TArgs, TState>): TableBufferingVgiFunction
source
export interface TableBufferingBindParams<TArgs = Record<string, any>>

Fields

argsTArgs
bindCallBindRequest
settingsRecord<string, any>
secretsRecord<string, Record<string, any>>
resolvedSecretsProvidedboolean

True on the second bind pass, after the connector resolved the secrets requested via the onBind lookupSecret* return fields on the first pass.

source
export interface TableBufferingConfig<
TArgs = Record<string, any>,
TState = any,
>

Fields

namestring
process( batch: VgiBatch, params: TableBufferingParams<TArgs>, ) => Uint8Array | Promise<Uint8Array>

Sink: ingest one batch, return an opaque state_id.

combine( stateIds: Uint8Array[], params: TableBufferingParams<TArgs>, ) => Uint8Array[] | Promise<Uint8Array[]>

Combine: group/merge state_ids, return finalize_state_ids.

finalize( params: TableBufferingParams<TArgs>, finalizeStateId: Uint8Array, state: TState, out: OutputCollector, ) => void | Promise<void>

Source tick: emit one batch via out.emit / signal EOS via out.finish.

descriptionstringoptional
argsRecord<string, VgiDataType>optional
namedArgsRecord<string, VgiDataType>optional
argDefaultsRecord<string, any>optional
onBind(params: TableBufferingBindParams<TArgs>) => | { outputSchema: VgiSchema; opaqueData?: Uint8Array; lookupSecretTypes?: string[]; lookupScopes?: string[]; lookupNames?: string[]; } | Promise<{ outputSchema: VgiSchema; opaqueData?: Uint8Array; lookupSecretTypes?: string[]; lookupScopes?: string[]; lookupNames?: string[]; }>optional

Bind: default passes through input schema. May be async. Return non-empty lookupSecret* lists on the first pass to request DuckDB secrets; the connector then re-binds with resolvedSecretsProvided and the resolved values are available in the process/finalize params’ secrets.

initialFinalizeState( finalizeStateId: Uint8Array, params: TableBufferingParams<TArgs>, ) => TState | Promise<TState>optional

Build the initial per-tick finalize state for a finalize_state_id.

cardinality( params: TableBufferingBindParams<TArgs>, ) => TableCardinality | Promise<TableCardinality>optional
projectionPushdownbooleanoptional
filterPushdownbooleanoptional
autoApplyFiltersbooleanoptional
sinkOrderDependentbooleanoptional

Force ParallelSink=false in the C++ operator (single-thread ingest).

sourceOrderDependentbooleanoptional

Force serial Source drain in finalize_queue order.

requiresInputBatchIndexbooleanoptional

Thread DuckDB’s per-chunk batch_index into every process() call.

stabilityFunctionStabilityoptional
examplesFunctionExample[]optional
categoriesstring[]optional
tagsRecord<string, string>optional
maxWorkersnumberoptional
requiredSettingsstring[]optional
requiredSecretsstring[]optional
source
export interface TableBufferingParams<TArgs = Record<string, any>>

Fields

argsTArgs
initCallInitRequest
outputSchemaVgiSchema
settingsRecord<string, any>
secretsRecord<string, Record<string, any>>
storageBoundStorage

Shared cross-process storage scoped by execution_id.

executionIdUint8Array

Stable across coordinator + secondary workers for one DuckDB execution.

attachIdUint8Array

Catalog attach identity (plaintext bytes).

transactionIdUint8Array | null

Hex/raw VGI transaction id, or null.

function_namestring
batchIndexnumber | null

Per-chunk batch_index when requiresInputBatchIndex=true; else null.

clientLog(level: string, message: string) => void

In-band log sink for the unary process()/combine() RPCs.

source
export interface TableBufferingVgiFunction extends VgiFunction

Description

A table_buffering VgiFunction also carries its callback config so the protocol’s unary handlers (process/combine/destructor) can reach it.

Fields

bufferingConfigTableBufferingConfig<any, any>
bufferingExtractArgs(request: BindRequest) => any