Buffering functions
On this page
Sink, combine, source — for output that depends on the whole input.
function defineTableBufferingFunction
Section titled “function defineTableBufferingFunction”export function defineTableBufferingFunction<TArgs = Record<string, any>,TState = any,>(config: TableBufferingConfig<TArgs, TState>): TableBufferingVgiFunctioninterface TableBufferingBindParams
Section titled “interface TableBufferingBindParams”export interface TableBufferingBindParams<TArgs = Record<string, any>>Fields
argsTArgsbindCallBindRequestsettingsRecord<string, any>secretsRecord<string, Record<string, any>>resolvedSecretsProvidedbooleanTrue on the second bind pass, after the connector resolved the secrets requested via the onBind
lookupSecret*return fields on the first pass.
interface TableBufferingConfig
Section titled “interface TableBufferingConfig”export interface TableBufferingConfig<TArgs = Record<string, any>,TState = any,>Fields
namestringprocess( 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.
descriptionstringoptionalargsRecord<string, VgiDataType>optionalnamedArgsRecord<string, VgiDataType>optionalargDefaultsRecord<string, any>optionalonBind(params: TableBufferingBindParams<TArgs>) => | { outputSchema: VgiSchema; opaqueData?: Uint8Array; lookupSecretTypes?: string[]; lookupScopes?: string[]; lookupNames?: string[]; } | Promise<{ outputSchema: VgiSchema; opaqueData?: Uint8Array; lookupSecretTypes?: string[]; lookupScopes?: string[]; lookupNames?: string[]; }>optionalBind: 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 withresolvedSecretsProvidedand the resolved values are available in the process/finalize params’secrets.initialFinalizeState( finalizeStateId: Uint8Array, params: TableBufferingParams<TArgs>, ) => TState | Promise<TState>optionalBuild the initial per-tick finalize state for a finalize_state_id.
cardinality( params: TableBufferingBindParams<TArgs>, ) => TableCardinality | Promise<TableCardinality>optionalprojectionPushdownbooleanoptionalfilterPushdownbooleanoptionalautoApplyFiltersbooleanoptionalsinkOrderDependentbooleanoptionalForce ParallelSink=false in the C++ operator (single-thread ingest).
sourceOrderDependentbooleanoptionalForce serial Source drain in finalize_queue order.
requiresInputBatchIndexbooleanoptionalThread DuckDB’s per-chunk batch_index into every process() call.
stabilityFunctionStabilityoptionalexamplesFunctionExample[]optionalcategoriesstring[]optionaltagsRecord<string, string>optionalmaxWorkersnumberoptionalrequiredSettingsstring[]optionalrequiredSecretsstring[]optional
interface TableBufferingParams
Section titled “interface TableBufferingParams”export interface TableBufferingParams<TArgs = Record<string, any>>Fields
argsTArgsinitCallInitRequestoutputSchemaVgiSchemasettingsRecord<string, any>secretsRecord<string, Record<string, any>>storageBoundStorageShared cross-process storage scoped by execution_id.
executionIdUint8ArrayStable across coordinator + secondary workers for one DuckDB execution.
attachIdUint8ArrayCatalog attach identity (plaintext bytes).
transactionIdUint8Array | nullHex/raw VGI transaction id, or null.
function_namestringbatchIndexnumber | nullPer-chunk batch_index when requiresInputBatchIndex=true; else null.
clientLog(level: string, message: string) => voidIn-band log sink for the unary process()/combine() RPCs.
interface TableBufferingVgiFunction
Section titled “interface TableBufferingVgiFunction”export interface TableBufferingVgiFunction extends VgiFunctionDescription
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