Load Confluent Kafka data to DuckDB
Build a Confluent Kafka to DuckDB pipeline with your coding agent. One prompt scaffolds it with the dltHub AI harness, plus the Confluent Kafka API base URL, auth, endpoints, and incremental loading.
Confluent Cloud provides REST APIs for managing Kafka resources, clusters, and security entities. Everything needed to build a working Confluent Kafka → 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 Kafka 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 Kafka 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 Kafka 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 Kafka API at a glance
| Base URL | For Confluent Cloud Kafka REST, the base URL is cluster-specific, e.g., 'https://<cluster-id>.<region>.<cloud>.confluent.cloud'. For general Confluent Cloud APIs, the base URL is 'https://api.confluent.cloud'. |
| Example endpoint | GET v3/clusters |
| Records found at | data |
| Authentication | all requests require an Authorization header using HTTP Basic authentication or a Bearer token for OAuth/STS flows — sent in the Authorization header, prefixed Bearer |
| Also required | Confluent-Identity-Pool-Id |
| Pagination | Cursor-based via page_token, page size via page_size (default 10, max 100). The page_token parameter is only valid on the first request; subsequent requests use the metadata.next URL for navigation. The Confluent Cloud API reference explicitly notes that Kafka V3 API list operations do not support pagination. |
| API reference | https://docs.confluent.io/cloud/current/kafka-rest/krest-qs.html |
These values come from the Confluent Kafka API reference — the authoritative source if anything here looks out of date.
How do I authenticate with the Confluent Kafka API?
Authentication is typically handled via an Authorization header using HTTP Basic authentication. The credentials consist of the API key ID and secret, colon-separated and base64-encoded, in the format 'Basic <base64_encoded_key_and_secret>'.
1. Get your credentials
- Navigate to the Confluent Cloud Console (https://confluent.cloud). 2. Expand the sidebar menu and select API keys, or navigate directly to the API keys page via your account or service account profile page. 3. Click Add API key. 4. Select the desired resource scope (e.g., Global for management APIs or specific Kafka cluster for data APIs). 5. Follow the prompts to name and describe the key. 6. Once the key and secret are generated, download or copy them immediately, as the secret cannot be retrieved after the creation dialog is closed. 7. Mark the checkbox to confirm you have saved the credentials to proceed.
2. Add them to .dlt/secrets.toml
[sources.confluent_kafka_source] confluent_api_key_id = "YOUR_API_KEY_ID" confluent_api_secret = "YOUR_API_SECRET" # Base64 encoded 'api_key_id:api_secret' for Authorization header confluent_auth_header = "Basic [BASE64_ENCODED_CREDENTIALS]"
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 Kafka data can I load into DuckDB?
These are the Confluent Kafka endpoints dlt can load into DuckDB:
| Resource | Endpoint | Method | Data selector | Description |
|---|---|---|---|---|
| clusters | /v3/clusters | GET | data | Lists all known Kafka clusters. |
| topics | /v3/clusters/{cluster_id}/topics | GET | data | Lists all topics for a specific cluster. |
| cluster | /v3/clusters/{cluster_id} | GET | Gets a specific cluster. | |
| broker_configs | /v3/clusters/{cluster_id}/broker-configs | GET | data | Lists broker configurations. |
| consumer_groups | /v3/clusters/{cluster_id}/consumer-groups | GET | data | Lists consumer groups for a cluster. |
How do I load only new Confluent Kafka records?
The Confluent Kafka API reference does not document a timestamp or sequence field for these endpoints, so there is nothing to advertise here as verified. Pick a field from the endpoints table above that increases with every write, then set it as the cursor_path.
{"name": "clusters", "endpoint": { "path": "v3/clusters", # Replace with a field that increases on every write. "incremental": {"cursor_path": "REPLACE_ME", "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 Kafka pipeline look like?
A standard dlt REST API pipeline — the same code you would write by hand, loading /kafka/v3/clusters and /subjects from the Confluent Kafka API into DuckDB:
import dlt from dlt.sources.rest_api import RESTAPIConfig, rest_api_resources @dlt.source def confluent_kafka_source(api_key=dlt.secrets.value): config: RESTAPIConfig = { "client": { "base_url": "For Confluent Cloud Kafka REST, the base URL is cluster-specific, e.g., 'https://<cluster-id>.<region>.<cloud>.confluent.cloud'. For general Confluent Cloud APIs, the base URL is 'https://api.confluent.cloud'.", "auth": {"type": "http_basic", "username": "REPLACE_ME", "password": api_key}, }, "resources": [ {"name": "clusters", "endpoint": {"path": "v3/clusters", "data_selector": "data"}}, {"name": "topics", "endpoint": {"path": "v3/clusters/{cluster_id}/topics", "data_selector": "data"}} ], } yield from rest_api_resources(config) def load_confluent_kafka_to_duckdb() -> None: pipeline = dlt.pipeline( pipeline_name="confluent_kafka_pipeline", destination="duckdb", dataset_name="confluent_kafka_data", ) load_info = pipeline.run(confluent_kafka_source()) print(load_info) if __name__ == "__main__": load_confluent_kafka_to_duckdb()
Run it with python confluent_kafka_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 Kafka 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_kafka_pipeline").dataset() df = data.clusters.df() print(df.head())
SQL:
SELECT * FROM confluent_kafka_data.clusters LIMIT 10;
See querying your data with dataset and exploring it in marimo notebooks.
How do I deploy the Confluent Kafka 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 Kafka 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 Kafka 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 Kafka to DuckDB?
Request dlt skills, commands, AGENT.md files, and AI-native context.