New to Kestra?
Use blueprints to kickstart your first workflows.
Compact raw OTel GenAI traces with in-memory DuckDB into Parquet and stream to BigQuery with token anomaly alerts.
Diagram unavailable
We could not build the topology for this blueprint. The flow itself is valid, use the YAML on the left to run it.
Collect and compact raw OpenTelemetry JSON traces into columnar Snappy Parquet using in-memory DuckDB, ingest into a day-partitioned and clustered Google BigQuery lakehouse table, and evaluate GenAI token consumption against anomaly thresholds.
Observability pipelines for Generative AI applications frequently generate massive volumes of detailed trace spans recording model identifiers, prompt lengths, completion tokens, latency, and error states. Ingesting raw JSON lines directly into analytical data warehouses incurs significant storage overhead and costly scan operations. This blueprint solves this by employing an embedded DuckDB engine to compact raw JSON telemetry into optimized Snappy-compressed Parquet files in memory before loading them into BigQuery with day-partitioning on timestamp and clustering on (service_name, model).
batch_telemetry_webhook trigger (io.kestra.plugin.core.trigger.Webhook) allows OpenTelemetry collectors, FluentBit, or forwarder proxies to POST trace batches via HTTP endpoint using authentication key otel-trace-ingest.hourly_compaction_schedule trigger (io.kestra.plugin.core.trigger.Schedule) enables scheduled recurring compaction every hour (shipped disabled by default).ingest_otel_traces):gen_ai.request.model, gen_ai.usage.input_tokens, gen_ai.usage.output_tokens, gen_ai.response.finish_reasons, duration_ms, service_name, status_code).compact_traces_to_parquet):io.kestra.plugin.duckdb.Query to read the raw newline-delimited JSON traces using read_json_auto().COPY ... TO ... (FORMAT PARQUET, CODEC 'SNAPPY') statement directly in memory, exporting an optimized Parquet file to Kestra internal storage.load_parquet_to_bigquery):io.kestra.plugin.gcp.bigquery.Load to append the Snappy Parquet file into the designated BigQuery table.timestamp column (timePartitioningField: "timestamp", timePartitioningType: "DAY") and clusters by ["service_name", "model"] to minimize downstream query scan costs.evaluate_anomalies):io.kestra.plugin.core.flow.If task evaluates whether the maximum token consumption exceeded token_anomaly_threshold or if errors occurred in the batch.send_slack_alert is true, dispatches a Slack block kit message via io.kestra.plugin.notifications.slack.SlackIncomingWebhook summarizing the incident.errors):Building telemetry ingestion with custom microservices or standing up Spark/Flink clusters incurs significant operational maintenance and infrastructure costs. Kestra provides a serverless, declarative approach that couples event-driven webhooks, scheduled micro-batching, in-memory embedded query engines like DuckDB, and native BigQuery loader tasks with enterprise-grade error handling and secret governance.
createDisposition: CREATE_IF_NEEDED).GCP_SERVICE_ACCOUNT_JSON: Service account private key JSON with roles/bigquery.dataEditor and roles/bigquery.jobUser IAM roles.SLACK_WEBHOOK_URL: Slack Incoming Webhook URL for alert delivery (required if send_slack_alert is set to true).| Name | Type | Default | Description |
|---|---|---|---|
gcs_bucket |
STRING | company-telemetry-lake |
Google Cloud Storage bucket containing raw OpenTelemetry JSON traces. |
gcs_prefix |
STRING | traces/raw/ |
Object prefix path for incoming raw trace files. |
bq_project |
STRING | company-analytics |
Google Cloud project ID hosting BigQuery. |
bq_dataset |
STRING | observability |
Target BigQuery dataset name. |
bq_table |
STRING | otel_genai_spans |
Target BigQuery table for compacted spans. |
token_anomaly_threshold |
INT | 50000 |
Token consumption ceiling above which alerts are raised. |
send_slack_alert |
BOOLEAN | false |
Enable or disable Slack alerting. |
GCP_SERVICE_ACCOUNT_JSON secret in your Kestra namespace with BigQuery permissions.SLACK_WEBHOOK_URL in your namespace secrets if alerting is desired.disabled: false on the hourly_compaction_schedule trigger.otel-collector-config.yaml pointing to the Kestra webhook URL using header key: otel-trace-ingest.io.kestra.plugin.dbt.cli.DbtCLI) to transform raw spans into hourly cost-per-service aggregations.