No logo available for Confluent Schema Registry to DuckDB connector icon

Load Confluent Schema Registry data to DuckDB

Build a Confluent Schema Registry to DuckDB pipeline with your coding agent. One prompt scaffolds it with the dltHub AI harness, plus the Confluent Schema Registry API base URL, auth, endpoints, and incremental loading.

SourceConfluent Schema RegistryConfluent Schema Registry API DocumentationDestinationDuckDBIn-process analytical database. The default local destination for dlt pipelines.

Confluent Schema Registry is a RESTful service for registering, retrieving, and managing schemas for data streams. Everything needed to build a working Confluent Schema Registry → DuckDB pipeline is on this page: the API's base URL, authentication, endpoints, pagination and incremental field — plus a prompt that hands the whole job to your coding agent.


Build your Confluent Schema Registry to DuckDB pipeline

Paste this prompt into Claude, Codex, or Cursor. The agent does the rest.

Prompt
Run uvx dlthub-init@latest to build a pipeline from Confluent Schema Registry to DuckDB and run it on dltHub

That scaffolds a dltHub workspace and installs the dltHub AI harness — the project rules, the secrets-management skill, and the dlt MCP server your agent needs to work safely. From there it reads the Confluent Schema Registry API, proposes the endpoints to load, then writes, runs and validates the pipeline while you review rather than type. Credentials are inspected through MCP tools, so your agent never reads secrets.toml itself. How the LLM-native workflow works →

Prefer to write it yourself? Every fact the agent uses is below.


Confluent Schema Registry API at a glance

Base URLhttp://localhost:8081
Example endpointGET schemas
Authenticationsupports HTTP Basic authentication or OAuth (Bearer token) — sent in the Authorization header, prefixed Bearer
Also requiredConfluent-Identity-Pool-Id, target-sr-cluster
PaginationOffset-based via offset, page size via limit (default -1). The API uses offset-based pagination via the 'offset' and 'limit' query parameters. 'limit' is an integer representing the page size, where values less than 0 are ignored or interpreted as 'no limit'. The 'offset' parameter represents the pagination offset.
Incremental fieldoffset
API referencehttps://docs.confluent.io/platform/current/schema-registry/develop/api.html

These values come from the Confluent Schema Registry API reference — the authoritative source if anything here looks out of date.


How do I authenticate with the Confluent Schema Registry API?

The API supports HTTP Basic authentication (passing username and password) or OAuth (passing a Bearer token in the Authorization header). For Basic auth, the credentials are typically passed as a base64-encoded 'Authorization: Basic ' header or via the --user flag in tools like curl.

1. Get your credentials

To obtain API credentials for Confluent Cloud Schema Registry: 1. Log in to the Confluent Cloud Console. 2. Navigate to the specific environment where your Schema Registry is located. 3. Select 'API Keys' from the left-hand navigation menu. 4. Click 'Add key' or 'Create key'. 5. Select 'Schema Registry' as the resource type. 6. Choose 'Service account' (recommended for production) or 'My account'. 7. Follow the prompts to generate the API key and secret. Ensure you save the secret immediately, as it cannot be retrieved later.

2. Add them to .dlt/secrets.toml

[sources.confluent_schema_registry_source] schema_registry_api_key = "your_key_here" schema_registry_api_secret = "your_secret_here" schema_registry_url = "https://your-schema-registry-endpoint.confluent.cloud"

dlt reads this file automatically at runtime. With the harness, the setup-secrets skill prompts you for the values and never handles the raw credential in chat. For production, see setting up credentials with dlt.


What Confluent Schema Registry data can I load into DuckDB?

These are the Confluent Schema Registry endpoints dlt can load into DuckDB:

ResourceEndpointMethodData selectorDescription
subjects/subjectsGETLists all registered subjects
schemas/schemasGETLists all schemas
schema_types/schemas/typesGETLists all supported schema types
mode/modeGETGets global configuration mode
config/configGETGets global cluster-level configuration

How do I load only new Confluent Schema Registry records?

Confluent Schema Registry exposes offset on schemas, so dlt can request only the records that changed since the last run. Set it as the cursor_path and dlt tracks the high-water mark for you between runs.

{"name": "schemas", "endpoint": { "path": "schemas", "incremental": {"cursor_path": "offset", "initial_value": "2024-01-01T00:00:00Z"}, }}

On the first run dlt loads everything from initial_value; on every run after that it requests only what changed and appends with write_disposition="merge" if you set a primary key. See incremental loading.


What does the generated Confluent Schema Registry pipeline look like?

A standard dlt REST API pipeline — the same code you would write by hand, loading subjects and schemas from the Confluent Schema Registry API into DuckDB:

import dlt from dlt.sources.rest_api import RESTAPIConfig, rest_api_resources @dlt.source def confluent_schema_registry_source(api_key=dlt.secrets.value): config: RESTAPIConfig = { "client": { "base_url": "http://localhost:8081", "auth": {"type": "bearer", "token": api_key}, }, "resources": [ {"name": "schemas", "endpoint": {"path": "schemas"}}, {"name": "subjects", "endpoint": {"path": "subjects"}} ], } yield from rest_api_resources(config) def load_confluent_schema_registry_to_duckdb() -> None: pipeline = dlt.pipeline( pipeline_name="confluent_schema_registry_pipeline", destination="duckdb", dataset_name="confluent_schema_registry_data", ) load_info = pipeline.run(confluent_schema_registry_source()) print(load_info) if __name__ == "__main__": load_confluent_schema_registry_to_duckdb()

Run it with python confluent_schema_registry_pipeline.py. The agent iterates on this until it loads cleanly — you review and approve, rather than write it from scratch.


How do I query Confluent Schema Registry data in DuckDB?

dlt creates one table per resource. Query the loaded data with Python or SQL — or ask your agent to, through the MCP server's execute_sql_query tool.

Python (pandas DataFrame):

import dlt data = dlt.pipeline("confluent_schema_registry_pipeline").dataset() df = data.schemas.df() print(df.head())

SQL:

SELECT * FROM confluent_schema_registry_data.schemas LIMIT 10;

See querying your data with dataset and exploring it in marimo notebooks.


How do I deploy the Confluent Schema Registry to DuckDB pipeline in production?

The pipeline runs locally, which is ideal for prototyping and one-off analysis. When you need it on a schedule, monitored on every load, and shared with your team, deploy the same dlt code on the dltHub platform — no infrastructure to maintain. The prompt above already ends with "run it on dltHub", so your agent can take it there directly.

  • Deploy & schedule — run the pipeline as a managed job with automatic retries.
  • Monitor — observable job queues, alerting, and load metrics for every run.
  • Transform — promote raw Confluent Schema Registry loads into governed, documented models.
  • Visualize & share — explore data in notebooks and publish live dashboards instead of static screenshots.

Book a demo →


What other destinations can I load Confluent Schema Registry data to?

dlt loads into any of these — only the destination argument changes:

DestinationExample value
PostgreSQL"postgres"
BigQuery"bigquery"
Snowflake"snowflake"
Redshift"redshift"
Databricks"databricks"
Filesystem (S3, GCS, Azure)"filesystem"

Set dlt.pipeline(destination="snowflake") and add credentials in .dlt/secrets.toml. On the dltHub platform the same pipeline runs against a managed Iceberg lakehouse. See the full destinations list.


Next steps

Was this page helpful?

Community Hub

Need more dlt context for Confluent Schema Registry to DuckDB?

Request dlt skills, commands, AGENT.md files, and AI-native context.