cdc streaming ingestion

    1

    Implement real-time RAG ingestion using CDC, Kafka, and idempotent vector database upserts.

    Free

    1 installSecurity scanned

    Works with the AI tools you already use

    Claude CodeClaude CodeCursorCursorCodex CLICodex CLIGitHub CopilotGitHub CopilotGemini CLIGemini CLI+20 more

    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

    ComponentImplementation Detail
    IdentityPoint IDs generated as {doc_id}-{chunk_index} for stable updates
    OrderingKafka partition key pinned to Primary Key to ensure sequential processing
    ConsistencyLSN (Log Sequence Number) stored in payload to skip stale/replay events
    DeletesExplicit 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_slots to ensure the slot doesn't grow indefinitely.
    • Initialize the Qdrant collection with a payload index on doc_id for 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.

    ~30 seconds
    1. 1

      Download the ZIP

      Free skills download straight away. Paid skills unlock right after purchase.

    2. 2

      Unzip into your skills folder

      Every agent reads skills from one folder on your machine. Drop the unzipped folder in there.

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

    Listed12 days ago

    What's inside

    Frequently Asked Questions