cdc streaming ingestion
Implement real-time RAG ingestion using CDC, Kafka, and idempotent vector database upserts.
Free
Works with the AI tools you already use
cdc streaming ingestion
Example session with this skill installed
Set up a Postgres CDC pipeline using Debezium and Kafka to sync my 'articles' table to a Qdrant vector store with exactly-once semantics.
- Read your context and instructions
- Compiled the cdc streaming ingestion
I have configured the CDC pipeline for real-time document ingestion. This setup includes the Postgres publication, the Debezium connector configuration, and a Python consumer that handles idempotent upserts to Qdrant.
Real-Time RAG Ingestion Pipeline
1. Postgres Publication and Slot
-- Enable logical replication and create publication
ALTER SYSTEM SET wal_level = 'logical';
SELECT pg_reload_conf();
CREATE PUBLICATION rag_pub FOR TABLE articles;
SELECT pg_create_logical_replication_slot('debezium_rag', 'pgoutput');
2. Debezium Connector Configuration
{
"name": "postgres-articles-connector",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"database.hostname": "db-host",
"database.port": "5432",
"database.dbname": "app_db",
"plugin.name": "pgoutput",
"slot.name": "debezium_rag",
"publication.name": "rag_pub",
"tombstones.on.delete": "true",
"key.converter": "org.apache.kafka.connect.json.JsonConverter",
"value.converter": "org.apache.kafka.connect.json.JsonConverter"
}
}
3. Idempotent Python Consumer
| Component | Implementation Detail |
|---|---|
| Identity | Point IDs generated as {doc_id}-{chunk_index} for stable updates |
| Ordering | Kafka partition key pinned to Primary Key to ensure sequential processing |
| Consistency | LSN (Log Sequence Number) stored in payload to skip stale/replay events |
| Deletes | Explicit handling of op: d and Kafka tombstone (null) messages |
async def idempotent_upsert(qdrant, point, event_lsn):
# Retrieve existing LSN to prevent out-of-order overwrites
existing = await qdrant.retrieve(collection_name="kb", ids=[point.id])
if existing and existing[0].payload.get("lsn", 0) >= event_lsn:
return
await qdrant.upsert(collection_name="kb", points=[point])
Next steps
- Deploy the Debezium connector to your Kafka Connect cluster.
- Monitor
pg_replication_slotsto ensure the slot doesn't grow indefinitely. - Initialize the Qdrant collection with a payload index on
doc_idfor efficient deletes.
Connects securely to your tools. The creator never sees your data.
About this skill
The problem
Scheduled batch ingestion leaves your vector database stale for hours. When high-throughput source data changes rapidly or compliance requires immediate deletion, batch processing cannot keep up with the freshness demands of real-time RAG applications.
What it does
- Configures Debezium and Postgres logical replication to capture row-level changes.
- Implements Kafka-based stream processing for embedding and upserting documents.
- Enforces exactly-once semantics through idempotent vector DB writes and LSN tracking.
- Handles document updates, deletes, and tombstones to ensure vector store consistency.
- Provides strategies for schema evolution and blue-green re-indexing during model upgrades.
Frameworks & tools
Postgres (WAL/Logical Replication), Debezium, Apache Kafka, Apache Flink, Python (aiokafka), Qdrant, and OpenAI Embeddings.
Why this beats prompting it yourself
Building streaming RAG requires solving complex distributed systems problems like late-arriving updates and offset management. This skill provides the specific infrastructure configurations and consumer logic needed to prevent data loss and WAL accumulation that generic prompts overlook.
Use cases
- Propagating sensitive data deletions to vector stores for GDPR compliance.
- Maintaining sub-10s freshness for RAG systems powered by high-traffic databases.
- Enriching document streams with temporal joins using Flink SQL.
- Scaling embedding pipelines for ingestion rates exceeding 100 writes per second.
Known limitations
Vector databases typically lack transactional parity with Kafka, requiring the idempotent upsert patterns provided here rather than native EOS.
How to install
Works the same in every agent - Claude, Cursor, Codex, Copilot and 20+ more.
- 1
Download the ZIP
Free skills download straight away. Paid skills unlock right after purchase.
- 2
Unzip into your skills folder
Every agent reads skills from one folder on your machine. Drop the unzipped folder in there.
- 3
Ask your agent to use it
Restart the agent if it was already running. It picks the skill up automatically - no config needed.
Skills folder by agent
Click the path to copy it. Create the folder if it does not exist yet.
Reviews
No reviews yet
Be one of the first to try it. Every listed skill passes our trust checks below.
Security scanned
Passed our 8-point scan before listing
1 install
Downloaded by developers to date
Free forever
No account required to browse
Trust & safety
Security scanned
Verified clean 12 days ago
- Free to download with an account