How Pathway Live Data Synchronization Handles Incremental Updates from Connected Data Sources
Pathway uses reactive computational graphs and dynamic tables to automatically propagate row-level deltas from source connectors through downstream transformations without polling or full reloads.
The pathwaycom/llm-app templates demonstrate a stream-processing engine designed for real-time LLM applications. Unlike batch-processing frameworks that require manual refresh cycles, Pathway's live data synchronization treats every connected source as a continuously evolving dataset, ensuring that embeddings, indices, and LLM contexts update instantly when underlying files change.
Core Architecture of Incremental Processing
Pathway represents every data source as a pw.Table—a dynamic, append-only collection that tracks insertions, modifications, and deletions. When connectors detect changes in external storage, they emit incremental update events that the runtime propagates through the computational graph.
Dynamic Tables and Source Connectors
Live source connectors watch external storage systems and translate file system events into table operations. In templates/unstructured_to_sql_on_the_fly/app.py, the pw.io.fs.read connector creates a table that monitors a directory for new or modified PDFs:
files = pw.io.fs.read(
data_dir, # e.g. "./data/quarterly_earnings"
format="binary",
) # ← automatically streams file-level changes
This connector leverages OS-level filesystem watchers to push row-level deltas into the graph. When a file appears or changes, Pathway inserts or updates the corresponding row without reloading existing data.
The Reactive Computational Graph
Pathway's runtime executes only the subgraph affected by changed rows. This incremental recomputation preserves previously computed results for unchanged data. In templates/drive_alert/app.py, when Google Drive files change, the system re-embeds only the modified documents and updates the KNNIndex without rebuilding the entire vector store:
files = pw.io.gdrive.read(
object_id=object_id,
service_user_credentials_file=service_user_credentials_file,
refresh_interval=30, # fetch changes every 30 s
)
parser = UnstructuredParser()
documents = files.select(texts=parser(pw.this.data)).flatten(pw.this.texts)
# Embed each new/updated chunk and update the nearest-neighbor index
enriched = documents + documents.select(data=embedder(pw.this.texts))
index = KNNIndex(enriched.data, enriched, n_dimensions=embedding_dimension)
Handling State-Aware Transformations
Downstream operations maintain internal state to enable complex event processing and avoid redundant work.
Stateful Deduplication Operators
The pw.stateful.deduplicate operator tracks previous outputs to filter noise and emit alerts only on substantive changes. As implemented in templates/drive_alert/app.py, this prevents duplicate LLM calls when file modifications do not alter the semantic content:
responses = prompt.select(
pw.this.query_id,
pw.this.query,
pw.this.alert_enabled,
response=model(prompt_chat_single_qa(pw.this.prompt)),
)
# Only keep changes that differ from the previous answer
dedup = pw.stateful.deduplicate(
responses,
col=responses.response,
acceptor=acceptor, # LLM decides if the change is significant
instance=responses.query_id,
)
This stateful operator compares new LLM responses against historical results for each query_id, guaranteeing that alerting logic triggers exclusively on meaningful deltas.
Persistence and Fault Tolerance
For production deployments, Pathway supports optional persistence layers that store intermediate table states. The templates/document_indexing/app.py implementation demonstrates how to configure backend storage (filesystem, Redis, or PostgreSQL) to enable resumable incremental updates after restarts. This ensures that the pipeline maintains its synchronization position with source systems even through process failures, preventing data loss or reprocessing storms.
Summary
- Reactive Graph Execution: Pathway's runtime automatically propagates changes from source connectors through transformations without explicit polling loops.
- Row-Level Deltas: Connectors like
pw.io.fs.readandpw.io.gdrive.reademit granular update events (insert/update/delete) rather than full dataset snapshots. - Selective Recomputation: Only operators dependent on changed rows re-execute, minimizing computational overhead for large datasets.
- Stateful Processing: Operators such as
pw.stateful.deduplicatemaintain historical context to enable change-aware alerting and avoid duplicate processing. - Fault Tolerance: Optional persistence layers in
document_indexing/app.pyallow incremental pipelines to resume state after interruptions.
Frequently Asked Questions
How does Pathway detect changes in local file systems without polling?
Pathway's pw.io.fs.read connector utilizes operating-system-specific filesystem watchers (inotify on Linux, FSEvents on macOS, ReadDirectoryChangesW on Windows) to receive push notifications when files are created, modified, or deleted. As shown in templates/unstructured_to_sql_on_the_fly/app.py, this creates a truly live table that reacts to events within milliseconds rather than scanning directories on a cron schedule.
What happens to the vector index when a source document is updated?
The KNNIndex and similar vector stores in Pathway receive incremental updates. When pw.io.gdrive.read detects a changed file in templates/drive_alert/app.py, the pipeline re-embeds only the affected document chunks and merges them into the existing index. The runtime tracks which embeddings derive from which source rows, enabling precise updates without full reindexing.
Can Pathway handle deletions and modifications, or only new files?
Pathway's dynamic tables support full change data capture (CDC), including deletions and modifications. Connectors emit tombstone records for deleted files and replacement rows for modifications. Downstream stateful operators like pw.stateful.deduplicate process these events to retract old answers and emit new ones, ensuring that the LLM application never references stale or removed content.
How does stateful deduplication prevent duplicate alerts?
The pw.stateful.deduplicate operator maintains a keyed state store (indexed by the instance parameter, typically a query or document ID) that compares incoming values against historically accepted ones using the specified acceptor function. In templates/drive_alert/app.py, this function uses an LLM judge to determine semantic equivalence; only rows that fail the equivalence test generate new downstream events, suppressing redundant notifications.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →