Distributed Fan-Out September 28, 2026 • 6 min read

Distributed Fan-Out & Multi-Sink Pipelines: Announcing the Apache SeaTunnel Output Connector

Modern enterprise Retrieval-Augmented Generation (RAG) architectures rarely write documents to a single destination. Knowledge chunks, metadata, Zero-Trust ACLs, and dense vector embeddings must simultaneously populate vector databases for semantic retrieval, lakehouses for compliance, search engines for keyword lookup, and data warehouses for audit logging. Today, we are thrilled to announce OpenCrawling's native Apache SeaTunnel Output Connector (oc-seatunnel-output-connector), integrating Apache SeaTunnel 2.3.13 and its high-throughput Zeta Engine for seamless distributed fan-out.

01. The Challenge: Multi-Destination Ingestion at Enterprise Scale

In traditional search architectures, crawler pipelines followed a point-to-point pattern: documents were extracted, chunked, embedded, and dumped into a single index. In modern enterprise AI environments, that model quickly breaks down:

Writing separate point-to-point sync jobs for each destination leads to duplicate crawler load, wasted embedding compute, and synchronization lag. The Apache SeaTunnel Output Connector resolves this by serving as a bridge between OpenCrawling's event-driven microservices and SeaTunnel's 100+ production-grade destination connectors.

02. Decoupled Architecture & Kafka Fan-Out

The integration leverages OpenCrawling's distributed microservice topology. The ingestion core extracts document text via Apache Tika, splits tokens, and passes chunks to oc-embedding-service for parallel vector computation. The resulting enriched payloads flow into Kafka:

Decoupled SeaTunnel Fan-Out Pipeline
Crawler
→
Kafka (Ingestion)
→
Tika Extractor
→
oc-embedding-service
→
SeaTunnel Writer Consumer
→
SeaTunnel Zeta Cluster
→
ClickHouse / Milvus / Iceberg

The connector consists of two key components:

  1. SeaTunnelOutputConnector (Control Plane): Synthesizes SeaTunnel HOCON Directed Acyclic Graphs (DAGs) and submits streaming jobs to the SeaTunnel Zeta cluster using its REST API v2.
  2. SeaTunnelStoreWriterConsumer (Data Plane): Subscribes to the Kafka opencrawling-embedded topic and streams Open Ingestion Standard (OIS) compliant records to an intermediate topic consumed by SeaTunnel's official Kafka source.

03. Automated HOCON DAG Synthesis

Writing complex HOCON configuration files for Apache SeaTunnel by hand can be error-prone. The connector includes SeaTunnelJobConfigBuilder, which compiles runtime configuration properties, execution parallelism, checkpoint intervals, Kafka broker credentials, catalog schemas, and chosen sink destinations into clean, validated HOCON definitions:

env {
  execution.parallelism = 4
  job.mode = "STREAMING"
  checkpoint.interval = 5000
}

source {
  Kafka {
    bootstrap.servers = "localhost:9092"
    topic = "opencrawling-embedded"
    consumer.group = "seatunnel-ingestion-consumer"
    result_table_name = "ois_embedded_chunks"
    format = "json"
    schema = {
      fields {
        id = "string"
        doc_id = "string"
        action = "string"
        uri = "string"
        text = "string"
        vector = "array<float>"
        acl = "array<string>"
        security_allowed_read = "array<string>"
        security_denied_read = "array<string>"
        security_inheritance = "boolean"
        last_modified = "string"
        metadata = "string"
      }
    }
  }
}

sink {
  Clickhouse {
    source_table_name = "ois_embedded_chunks"
    host = "localhost:8123"
    database = "default"
    table = "document_chunks"
  }
  Milvus {
    source_table_name = "ois_embedded_chunks"
    url = "http://localhost:19530"
    collection = "enterprise_kb"
  }
}

04. SeaTunnel Zeta REST API v2 Control Plane

The connector features a lightweight, robust HTTP/1.1 client—SeaTunnelRestClient—built on Java 25 HttpClient. It communicates directly with SeaTunnel Zeta master nodes (typically on port 8088 or 8080) to automate full pipeline lifecycles:

05. OIS Schema & Zero-Trust Security ACL Mapping

OpenCrawling's Open Ingestion Standard (OIS) guarantees that security permissions are never detached from content chunks. SeaTunnel catalog schemas carry the full security model across every downstream sink:

Field Type Description
id string Unique chunk identifier (e.g. doc-001_chunk_0)
doc_id string Parent document identifier for tracking multi-chunk lineage
vector array<float> Dense embedding vector (e.g. 1024 floats)
acl array<string> Zero-Trust Security SIDs (e.g. ["ROLE_USER", "group:engineers"])
security_allowed_read array<string> Explicit allowed reader identities for strict ACL evaluation
security_denied_read array<string> Explicit denied reader identities overriding read permissions
security_inheritance boolean Security inheritance flag preserved from repository hierarchy

06. Document Lifecycle CDC & Tombstone Deletes

When enterprise documents are deleted from source systems (such as Alfresco, SharePoint, or web repositories), re-embedding every chunk to maintain search index cleanliness is cost-prohibitive. OpenCrawling emits lightweight OIS Deletion Tombstones (action: DELETE).

The SeaTunnel connector maps these tombstones to SeaTunnel Change Data Capture (CDC) row kinds (-D / RowKind.DELETE). Downstream storage engines like ClickHouse, Iceberg, and Milvus automatically purge removed records without executing any redundant AI embedding calls.

💡

Zero-Cluster Footprint (Stock SeaTunnel Compatibility): OpenCrawling runs on Java 25 with preview features enabled (Virtual Threads & Structured Concurrency), while Apache SeaTunnel 2.3.13 engines typically run on Java 8, 11, or 17. By decoupling the connector over standard Kafka topics, users do not need to install custom JARs into /opt/seatunnel/connectors/, avoiding JVM classpath and class version conflicts.

07. Configuration Properties Reference

All properties can be configured in application.yml under the spring.opencrawling.output.seatunnel.* prefix or mapped to standard environment variables:

Property Default Description
rest-url http://localhost:8080 HTTP endpoint of the SeaTunnel Zeta master REST API
job-name opencrawling_ingestion_pipeline Identifier for the submitted SeaTunnel job
job-mode STREAMING Execution mode: STREAMING or BATCH
checkpoint-interval-ms 5000 State checkpoint interval in milliseconds
parallelism 4 Parallelism degree across Zeta workers
kafka-bootstrap-servers localhost:9092 Kafka bootstrap servers for the fan-out stream
kafka-topic opencrawling-embedded Topic containing embedded document chunks
target-sinks console Comma-separated list of target sinks (e.g. clickhouse,milvus)
auto-submit-job true Automatically submit the generated HOCON DAG on startup

08. Admin UI & Real-Time Diagnostics

The OpenCrawling Admin UI (oc-admin-ui) provides dedicated configuration forms for the SeaTunnel connector. Administrators can specify the master REST URL, job name, execution parallelism, checkpoint intervals, and target fan-out sinks. The Check Connection button invokes ConnectorCheckerService to dynamically ping the Zeta cluster via /overview.

09. Automated Decoupled Integration Testing

To ensure total reliability, OpenCrawling includes a dedicated multi-service integration test suite (scripts/test-seatunnel-decoupled.sh) backed by docker-compose-decoupled-with-seatunnel.yml:

# Run end-to-end decoupled integration test
./scripts/test-seatunnel-decoupled.sh

The script executes the full ingestion lifecycle:

  1. Spins up Apache SeaTunnel Zeta cluster (apache/seatunnel:2.3.13), Kafka, Zookeeper, Redis, Ollama (with mxbai-embed-large), and all OpenCrawling microservices.
  2. Submits test documents through the crawler, extracts text via Tika, and computes 1024-dimension embeddings.
  3. Serializes OIS records, streams to Kafka, and deploys the HOCON streaming pipeline to SeaTunnel Zeta.
  4. Verifies active execution via curl http://localhost:8088/running-jobs and validates MCP Server connectivity on port 8080.
  5. Executes unit tests verifying OIS Tombstone DELETE actions and tears down all containers cleanly.

Ready to Fan Out Your Enterprise RAG Pipelines?

Try the live interactive simulator, inspect the source code, or read our complete Wiki documentation.