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.
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.
PromptRunuvx dlthub-init@latestto 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 URL | http://localhost:8081 |
| Example endpoint | GET schemas |
| Authentication | supports HTTP Basic authentication or OAuth (Bearer token) — sent in the Authorization header, prefixed Bearer |
| Also required | Confluent-Identity-Pool-Id, target-sr-cluster |
| Pagination | Offset-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 field | offset |
| API reference | https://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:
| Resource | Endpoint | Method | Data selector | Description |
|---|---|---|---|---|
| subjects | /subjects | GET | Lists all registered subjects | |
| schemas | /schemas | GET | Lists all schemas | |
| schema_types | /schemas/types | GET | Lists all supported schema types | |
| mode | /mode | GET | Gets global configuration mode | |
| config | /config | GET | Gets 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.
What other destinations can I load Confluent Schema Registry data to?
dlt loads into any of these — only the destination argument changes:
| Destination | Example 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.