Webhook icon
Schedule icon
Script icon
Process icon
Query icon
Load icon
If icon
SlackIncomingWebhook icon
Log icon

OpenTelemetry GenAI & Trace Telemetry Pipeline to BigQuery via DuckDB

Compact raw OTel GenAI traces with in-memory DuckDB into Parquet and stream to BigQuery with token anomaly alerts.

Categories
AICloudData

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).

How it works

  1. Trigger Options:
    • The 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.
    • The hourly_compaction_schedule trigger (io.kestra.plugin.core.trigger.Schedule) enables scheduled recurring compaction every hour (shipped disabled by default).
  2. Ingestion & Validation (ingest_otel_traces):
    • A lightweight Python task generates or processes OpenTelemetry trace spans formatted in accordance with OpenTelemetry GenAI semantic conventions (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).
    • Calculates batch metrics including total spans, peak token consumption, and error counts, exposing them via Kestra outputs.
  3. In-Memory Compaction (compact_traces_to_parquet):
    • Uses io.kestra.plugin.duckdb.Query to read the raw newline-delimited JSON traces using read_json_auto().
    • Flattens nested JSON attributes and executes a COPY ... TO ... (FORMAT PARQUET, CODEC 'SNAPPY') statement directly in memory, exporting an optimized Parquet file to Kestra internal storage.
  4. Warehouse Load (load_parquet_to_bigquery):
    • Uses io.kestra.plugin.gcp.bigquery.Load to append the Snappy Parquet file into the designated BigQuery table.
    • Enforces day-partitioning on the timestamp column (timePartitioningField: "timestamp", timePartitioningType: "DAY") and clusters by ["service_name", "model"] to minimize downstream query scan costs.
  5. Flow Control & Anomaly Alerting (evaluate_anomalies):
    • An io.kestra.plugin.core.flow.If task evaluates whether the maximum token consumption exceeded token_anomaly_threshold or if errors occurred in the batch.
    • When triggered and send_slack_alert is true, dispatches a Slack block kit message via io.kestra.plugin.notifications.slack.SlackIncomingWebhook summarizing the incident.
  6. Error Handling (errors):
    • A flow-level error handler logs execution failures and delivers failure alerts to Slack if enabled.

What you get

  • Substantial cloud storage and query cost reductions by converting JSON to Snappy Parquet before loading.
  • Day-partitioned and clustered BigQuery tables optimized for observability dashboards, cost attribution, and SRE incident triage.
  • Automated detection and Slack alerting for runaway LLM prompt loops, large context spills, and API rate-limiting spikes.
  • Full end-to-end execution lineage and audit history in Kestra.

Who it's for

  • Platform Engineers and SREs monitoring production Generative AI deployments and microservice meshes.
  • FinOps teams tracking LLM token expenditure across models, teams, and services.
  • Data Engineers building lakehouse architectures on Google Cloud Platform without standing up permanent streaming infrastructure.

Why orchestrate this with Kestra

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.

Prerequisites

  • A Google Cloud Platform (GCP) project with the BigQuery API enabled.
  • A BigQuery dataset created (or service account permissions to allow createDisposition: CREATE_IF_NEEDED).
  • An active Slack Incoming Webhook URL (if notifications are enabled).

Secrets

  • 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).

Inputs

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.

Quick start runbook

  1. Configure the GCP_SERVICE_ACCOUNT_JSON secret in your Kestra namespace with BigQuery permissions.
  2. (Optional) Configure SLACK_WEBHOOK_URL in your namespace secrets if alerting is desired.
  3. Run the workflow manually from the Kestra UI to verify trace generation, DuckDB Parquet compaction, and BigQuery loading.
  4. To enable scheduled batch execution, edit the workflow and set disabled: false on the hourly_compaction_schedule trigger.
  5. To ingest live traces from OpenTelemetry Collectors, configure an HTTP exporter in your otel-collector-config.yaml pointing to the Kestra webhook URL using header key: otel-trace-ingest.

How to extend

  • Add a downstream dbt task (io.kestra.plugin.dbt.cli.DbtCLI) to transform raw spans into hourly cost-per-service aggregations.
  • Connect Grafana or Looker Studio to the clustered BigQuery table to visualize token latency P99 and token budget burn rates.
  • Incorporate data loss prevention or PII redaction rules within the DuckDB SQL compaction stage before writing Parquet files.

Links

See How

New to Kestra?

Use blueprints to kickstart your first workflows.