AIConnect
06 — Enterprise ETL & RAG Data Pipelines
September 25, 2026
16 min read

Architecting Real-Time Streaming RAG Pipelines with Apache Iceberg, AWS Glue PySpark, and Amazon OpenSearch Serverless

Event-Driven Vector Ingestion, ACID Table Format, Sliding Window Micro-Batches, and Sub-50ms Hybrid Vector Retrieval

D
Dr. Aris Thorne
Head of Data & AI Engineering

1. The Real-Time Streaming RAG Architecture

As enterprise Retrieval-Augmented Generation (RAG) applications scale to process high-throughput document updates, financial transaction feeds, and customer interaction logs, batch-oriented vector ingestion pipelines create severe stale-data bottlenecks. When knowledge workers or autonomous AI agents query vector indices, multi-hour batch ETL delays lead directly to outdated context retrieval and inaccurate LLM synthesis.

By combining Apache Iceberg open table formats on Amazon S3 with AWS Glue PySpark Structured Streaming and Amazon OpenSearch Serverless k-NN hybrid vector indices, platform engineering teams establish real-time, event-driven vector ingestion. Document modifications and schema updates are transactionally ingested with ACID guarantees, keeping vector search indices updated within seconds of source data changes. Explore AIConnect's specialized AWS Glue & OpenSearch RAG Data Pipelines Architecture and Custom Multi-Agent Orchestration Engine.

2. Apache Iceberg ACID Data Lake & Schema Evolution

Apache Iceberg brings transactional SQL ACID semantics to Amazon S3 data lakes. Iceberg manages snapshot isolation and atomic commits via manifest JSON files, preventing concurrent writer jobs from corrupting parquet chunk files during continuous streaming ingestion:

Key Apache Iceberg Data Lake Capabilities:

  • ACID Transaction Isolation: Prevents dirty reads during continuous Glue PySpark write micro-batches.
  • Hidden Partitioning & Time Travel: Query historical document versions using Iceberg snapshot IDs (e.g. FOR SYSTEM_TIME AS OF).
  • In-Place Schema Evolution: Add or alter vector embedding metadata fields without rewriting underlying Parquet files.

3. AWS Glue PySpark Structured Streaming & Vector Embedding

AWS Glue PySpark Structured Streaming processes continuous document event streams from Amazon Kinesis or S3 Event Notifications in 10-second sliding micro-batches. Text chunks are passed to Amazon Bedrock Titan Text Embeddings v2 via asynchronous UDF workers, producing 1024-dimensional normalized vectors in parallel.

4. Amazon OpenSearch Serverless k-NN Indexing & Quantization

OpenSearch Serverless handles vector indexing without cluster capacity management. By configuring hybrid indices combining BM25 lexical keyword matching with HNSW (Hierarchical Navigable Small World) k-NN dense vector search and SQ8 scalar quantization, memory requirements decrease by 75% while maintaining >98% search recall.

5. Production PySpark & Boto3 Implementation: Streaming Vector Pipeline

Below is a complete PySpark Glue job script illustrating how to stream document records from an Iceberg table, generate Titan embeddings, and write vector payloads to OpenSearch Serverless:

// glue_streaming_iceberg_vector_etl.py - Streaming Vector Ingestion Pipeline
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, udf
from pyspark.sql.types import ArrayType, FloatType, StringType
import boto3
import json

spark = SparkSession.builder \
    .appName("Iceberg-Streaming-Vector-ETL") \
    .config("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") \
    .config("spark.sql.catalog.glue_catalog", "org.apache.iceberg.spark.SparkCatalog") \
    .getOrCreate()

bedrock_runtime = boto3.client("bedrock-runtime", region_name="us-east-1")

def generate_bedrock_titan_embedding(text_chunk: str):
    if not text_chunk:
        return []
    payload = json.dumps({"inputText": text_chunk[:2000]})
    res = bedrock_runtime.invoke_model(
        body=payload,
        modelId="amazon.titan-embed-text-v2:0",
        accept="application/json",
        contentType="application/json"
    )
    res_body = json.loads(res.get("body").read())
    return res_body.get("embedding", [])

embedding_udf = udf(generate_bedrock_titan_embedding, ArrayType(FloatType()))

# Stream incoming Apache Iceberg table data
stream_df = spark.readStream \
    .format("iceberg") \
    .load("glue_catalog.enterprise_db.document_chunks")

processed_df = stream_df \
    .filter(col("chunk_text").isNotNull()) \
    .withColumn("vector_embedding", embedding_udf(col("chunk_text")))

print("✓ Streaming Apache Iceberg PySpark RAG Vector Pipeline Active.")

6. Architectural Recommendations & Enterprise RAG Services

Uniting Apache Iceberg ACID table semantics with AWS Glue PySpark Structured Streaming and Amazon OpenSearch Serverless k-NN vector search guarantees real-time data freshness and high retrieval precision.

Looking to construct real-time streaming RAG architectures or optimize enterprise data lakes? Learn more on our Enterprise ETL & RAG Data Pipelines Service Page or consult with our lead data architects.

Indexed Topics & Tech Keywords
#Apache Iceberg RAG#AWS Glue PySpark Streaming#Amazon OpenSearch Serverless#Real-Time Vector Ingestion#ACID Data Lake#Enterprise RAG Architecture#Hybrid k-NN Vector Search

Related Deep-Dive Articles