config

package
v1.2.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Sep 11, 2026 License: MIT Imports: 11 Imported by: 0

Documentation

Index

Constants

View Source
const (
	DefaultKafkaFetchMaxBytes          = 100 << 20
	DefaultKafkaFetchMaxPartitionBytes = 10 << 20
	DefaultKafkaFetchPrefetch          = 2
)

Defaults for KafkaFetch. The two byte values are what the source set before the block existed. The prefetch default is measured: the smallest depth within 5% of unbounded throughput on a 3M message backlog. See docs/superpowers/specs/2026-09-10-kafka-fetch-bound-design.md.

Variables

This section is empty.

Functions

func RenderTemplate

func RenderTemplate(path string, overrides map[string]string) ([]byte, error)

RenderTemplate renders a config through Jinja2, the templating the config spec is written in: every SQLFLOW_-prefixed environment variable is in scope, plus the settings vars, plus explicit overrides.

func RenderTemplateString added in v1.1.0

func RenderTemplateString(src []byte, overrides map[string]string) ([]byte, error)

RenderTemplateString renders config text. RenderTemplate is the same thing for a file, and calls this.

The text form exists because a validation job carries the config's content rather than a path: the pull-based control plane hands a job to an instance that shares no filesystem with the submitter (#178).

func SQLResultsCacheDir

func SQLResultsCacheDir() string

SQLResultsCacheDir mirrors sqlflow.settings.SQL_RESULTS_CACHE_DIR: where disk-backed handlers stage their batch and result files. It is the fallback for a handler that does not declare sql_results_cache_dir.

func SettingsVarNames added in v1.1.0

func SettingsVarNames() []string

SettingsVarNames are the variables the engine injects into every render.

Validation needs them separated from the rest: nobody supplied them, so reporting them as supplied-but-never-read would put two lines of noise on every config and bury the one name that matters.

func TemplateVars added in v1.1.0

func TemplateVars(overrides map[string]string) map[string]string

TemplateVars returns the context a render runs against, so validation can report what the config actually had available to it.

Types

type ClickhouseSink

type ClickhouseSink struct {
	DSN   string `yaml:"dsn"`
	Table string `yaml:"table"`
}

type Conf

type Conf struct {
	// Main pipeline configuration.
	Pipeline Pipeline `yaml:"pipeline"`
	// Predefined SQL tables used in the pipeline.
	Tables *Tables `yaml:"tables,omitempty"`
	// List of User-Defined Functions (UDFs) to be used in SQL queries.
	UDFs []UDF `yaml:"udfs,omitempty"`
	// List of SQL commands to execute before processing the pipeline.
	Commands []SQLCommand `yaml:"commands,omitempty"`
}

Conf is a whole pipeline file.

func Load

func Load(path string, overrides map[string]string) (*Conf, error)

Load is LoadRendered for callers that do not need the text.

func LoadRendered added in v1.2.0

func LoadRendered(path string, overrides map[string]string) (*Conf, []byte, error)

LoadRendered renders the file and parses it, returning both. The rendered text is what the TurboStats bundle hashes: two instances running the same file under different environments are running different configs, and the hash should say so.

type ConsoleSink

type ConsoleSink struct{}

type Error

type Error struct {
	Policy ErrorPolicy `yaml:"policy"`
}

Error

type ErrorPolicy

type ErrorPolicy string
const (
	PolicyRaise  ErrorPolicy = "RAISE"
	PolicyIgnore ErrorPolicy = "IGNORE"
)

type Handler

type Handler struct {
	Type               string `yaml:"type"`
	SQL                string `yaml:"sql"`
	SQLResultsCacheDir string `yaml:"sql_results_cache_dir,omitempty"`
	Table              string `yaml:"table,omitempty"`
}

Handler

type IcebergSink

type IcebergSink struct {
	CatalogName string `yaml:"catalog_name"`
	TableName   string `yaml:"table_name"`
}

Various Sink Configs

type KafkaFetch added in v1.1.0

type KafkaFetch struct {
	// Bytes one fetch may return per broker. Kafka's fetch.max.bytes.
	MaxBytes int `yaml:"max_bytes,omitempty" jsonschema:"minimum=1"`
	// Bytes one fetch may return per partition. Kafka's max.partition.fetch.bytes.
	MaxPartitionBytes int `yaml:"max_partition_bytes,omitempty" jsonschema:"minimum=1"`
	// Fetches held ahead of the pipeline.
	Prefetch int `yaml:"prefetch,omitempty" jsonschema:"minimum=1"`
}

KafkaFetch bounds the consumer's read-ahead. The source holds at most prefetch fetches in the channel the pipeline reads from, plus the one it is waiting to send, and franz-go holds one more per broker. In bytes of payload, the worst case is (prefetch + 2) x brokers x max_bytes and the typical case is (prefetch + 2) x partitions x max_partition_bytes.

func (*KafkaFetch) Resolved added in v1.1.0

func (f *KafkaFetch) Resolved() (KafkaFetch, error)

Resolved fills absent fields with their defaults and checks the bounds. A nil receiver is the absent block.

type KafkaSASL

type KafkaSASL struct {
	Mechanism string `yaml:"mechanism" jsonschema:"enum=PLAIN,enum=SCRAM-SHA-256,enum=SCRAM-SHA-512,enum=GSSAPI"`
	Username  string `yaml:"username"`
	Password  string `yaml:"password"`
}

KafkaSASL configures SASL authentication for a Kafka connection.

type KafkaSSL

type KafkaSSL struct {
	CALocation                      string `yaml:"ca_location,omitempty"`
	CertificateLocation             string `yaml:"certificate_location,omitempty"`
	KeyLocation                     string `yaml:"key_location,omitempty"`
	KeyPassword                     string `yaml:"key_password,omitempty"`
	EndpointIdentificationAlgorithm string `yaml:"endpoint_identification_algorithm,omitempty"`
}

KafkaSSL configures TLS material for a Kafka connection.

type KafkaSink

type KafkaSink struct {
	// List of Kafka brokers.
	Brokers []string `yaml:"brokers"`
	// Target Kafka topic.
	Topic            string     `yaml:"topic"`
	SecurityProtocol string     `yaml:"security_protocol,omitempty"`
	SSL              *KafkaSSL  `yaml:"ssl,omitempty"`
	SASL             *KafkaSASL `yaml:"sasl,omitempty"`
}

type KafkaSource

type KafkaSource struct {
	Brokers         []string `yaml:"brokers"`
	GroupID         string   `yaml:"group_id"`
	AutoOffsetReset string   `yaml:"auto_offset_reset" jsonschema:"enum=earliest,enum=latest"`
	Topics          []string `yaml:"topics"`

	SecurityProtocol string     `yaml:"security_protocol,omitempty" jsonschema:"enum=SASL_SSL,enum=SSL,enum=SASL_PLAINTEXT,enum=PLAINTEXT"`
	SSL              *KafkaSSL  `yaml:"ssl,omitempty"`
	SASL             *KafkaSASL `yaml:"sasl,omitempty"`

	// Bounds how far the consumer reads ahead of the pipeline. Omit the
	// block to accept the defaults, which bound a backlog replay to a few
	// fetches rather than the backlog.
	Fetch *KafkaFetch `yaml:"fetch,omitempty"`
}

Source Types

type OnError

type OnError struct {
	// Defines how errors should be handled.
	Policy string `yaml:"policy" jsonschema:"enum=RAISE,enum=IGNORE,enum=DLQ"`
	// Dead-letter queue configuration. Failed messages will be routed to this
	// sink.
	DLQ *Sink `yaml:"dlq,omitempty"`
}

OnError configures what happens to a message or batch that fails. The dlq block is a full sink definition, used when policy is DLQ.

type Pipeline

type Pipeline struct {
	// Name of the pipeline.
	Name string `yaml:"name,omitempty"`
	// Description of the pipeline.
	Description string `yaml:"description,omitempty"`
	// Configuration for the data source.
	Source  Source  `yaml:"source"`
	Handler Handler `yaml:"handler"`
	Sink    Sink    `yaml:"sink"`
	// Messages accumulated before the handler runs. Larger batches trade
	// latency for throughput.
	BatchSize int `yaml:"batch_size,omitempty"`
	// Longest a partial batch waits before it is invoked anyway.
	FlushIntervalSeconds int `yaml:"flush_interval_seconds,omitempty"`
	// Where the pipeline keeps its DuckDB state. Absent means in-memory, and
	// state is lost on a crash.
	State *StateConf `yaml:"state,omitempty"`
	// Global error handling strategy for the pipeline.
	OnError *OnError `yaml:"on_error,omitempty"`
}

Pipeline is the source, the query, and the destination.

type SQLCommand

type SQLCommand struct {
	// Name of the command for reference.
	Name string `yaml:"name"`
	// SQL statements to execute.
	SQL string `yaml:"sql"`
}

SQLCommand is one statement run at startup, before any table is created.

type SQLCommandSink

type SQLCommandSink struct {
	SQL           string                   `yaml:"sql"`
	Substitutions []SQLCommandSubstitution `yaml:"substitutions,omitempty"`
}

type SQLCommandSubstitution

type SQLCommandSubstitution struct {
	Var  string `yaml:"var"`
	Type string `yaml:"type" jsonschema:"enum=uuid4"`
}

type Sink

type Sink struct {
	// Sink identifier.
	Type string `yaml:"type"`
	// Format settings (e.g., for Parquet).
	Format *SinkFormat `yaml:"format,omitempty"`
	// Kafka-specific sink configuration.
	Kafka *KafkaSink `yaml:"kafka,omitempty"`
	// Console output sink configuration.
	Console *ConsoleSink `yaml:"console,omitempty"`
	// SQL-command sink configuration.
	SQLCommand *SQLCommandSink `yaml:"sqlcommand,omitempty"`
	// Iceberg-specific sink configuration.
	Iceberg *IcebergSink `yaml:"iceberg,omitempty"`
	// ClickHouse-specific sink configuration.
	Clickhouse *ClickhouseSink `yaml:"clickhouse,omitempty"`
	// Bounds how long this sink keeps trying a destination that is not
	// answering. Omit to accept the defaults; set max_attempts to 1 to turn
	// retrying off. The kafka sink ignores this: franz-go already retries a
	// produce with its own backoff.
	Retry *SinkRetry `yaml:"retry,omitempty"`
}

Sink is where result rows go. One block per destination, and the type field selects which one the pipeline builds.

The same shape serves three places: a pipeline's sink, the dead-letter queue a failed record diverts to, and the sink a managed window collects into.

type SinkFormat

type SinkFormat struct {
	Type string `yaml:"type" jsonschema:"enum=parquet"`
}

SinkFormat

Closed value sets carry jsonschema enum tags. Go has no enum type, so without them the generated schema would accept any string where the hand-written one accepted one of a few, and validation would get weaker as a side effect of generating it.

type SinkRetry added in v1.0.5

type SinkRetry struct {
	// Total attempts, including the first. 1 disables retrying.
	MaxAttempts int `yaml:"max_attempts,omitempty"`
	// Wait before the first retry. Doubles each attempt.
	InitialBackoffMS int `yaml:"initial_backoff_ms,omitempty"`
	// Ceiling on the backoff.
	MaxBackoffMS int `yaml:"max_backoff_ms,omitempty"`
	// Bounds the whole ladder, not one attempt. Keep it below
	// pipeline.flush_interval_seconds: the retry runs inside the open state
	// transaction, whose clock the window depends on.
	DeadlineSeconds int `yaml:"deadline_seconds,omitempty"`
}

SinkRetry bounds how long a sink keeps trying a destination that is not answering. Omit the block to accept the defaults; set max_attempts to 1 to turn retrying off.

The Kafka sink ignores this. franz-go already retries a produce with its own backoff, and a second ladder on top of that one is worse than none.

type Source

type Source struct {
	Type      string           `yaml:"type"`
	Kafka     *KafkaSource     `yaml:"kafka,omitempty"`
	Websocket *WebsocketSource `yaml:"websocket,omitempty"`
	Webhook   *WebhookSource   `yaml:"webhook,omitempty"`
	Error     *Error           `yaml:"error,omitempty"`
}

Source

type StateConf added in v1.0.5

type StateConf struct {
	// File backing the pipeline's DuckDB database. Window state and Kafka
	// offsets are stored here and committed together.
	Path string `yaml:"path"`
}

StateConf points the pipeline's DuckDB at a file, so tables the handler writes -- window state above all -- survive a restart. Offsets are stored in the same database, which is what makes state and offsets recoverable together.

type TableManager

type TableManager struct {
	TumblingWindow *TumblingWindow `yaml:"tumbling_window"`
	Sink           Sink            `yaml:"sink"`
}

Table Manager

type TableSQL

type TableSQL struct {
	Name    string        `yaml:"name"`
	SQL     string        `yaml:"sql"`
	Manager *TableManager `yaml:"manager,omitempty"`
}

SQL Tables

type Tables

type Tables struct {
	// List of tables with their SQL definitions and management configurations.
	SQL []TableSQL `yaml:"sql"`
}

Tables holds the SQL tables the pipeline creates before it consumes.

type TumblingWindow

type TumblingWindow struct {
	CollectSQL string `yaml:"collect_closed_windows_sql"`
	DeleteSQL  string `yaml:"delete_closed_windows_sql"`
	// How often to collect closed windows. Optional; the manager applies its
	// own default when this is absent or not positive.
	PollIntervalSecs int `yaml:"poll_interval_seconds,omitempty"`
}

Tumbling Window Manager

type UDF

type UDF struct {
	FunctionName string `yaml:"function_name"`
	ImportPath   string `yaml:"import_path"`
}

UDF registers a user-defined function the handler SQL can call.

type WebhookHMAC

type WebhookHMAC struct {
	Header string `yaml:"header"`
	SigKey string `yaml:"sig_key"`
	Secret string `yaml:"secret"`
}

WebhookHMAC configures signature validation of incoming webhook bodies. SigKey names the digest the signature header is prefixed with, e.g. "sha256" for GitHub's "X-Hub-Signature-256: sha256=<hex>".

type WebhookSource

type WebhookSource struct {
	SignatureType string       `yaml:"signature_type,omitempty" jsonschema:"enum=hmac"`
	HMAC          *WebhookHMAC `yaml:"hmac,omitempty"`
}

type WebsocketSource

type WebsocketSource struct {
	URI string `yaml:"uri"`
}

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL