How to Implement Real-Time Alerting When Document Content Changes in Pathway
Pathway's drive_alert template provides a production-ready pattern for monitoring connected sources like Google Drive and automatically notifying users via Slack when document updates materially change query answers.
The pathwaycom/llm-app repository contains a complete reference implementation for real-time alerting when document content changes across connected data sources. This solution leverages Pathway's reactive streaming engine to continuously index files, detect meaningful content deviations using LLM-based comparison logic, and trigger notifications only when answers to subscribed queries actually change.
Ingest and Index Documents from Connected Sources
Poll Google Drive for File Updates
The pipeline begins by establishing a persistent connection to your document source. In templates/drive_alert/app.py, the pw.io.gdrive.read function polls a specific Drive folder every 30 seconds using the refresh_interval=30 parameter, emitting new rows whenever files are added or modified【L63-L67】. Each file is streamed as a dynamic table row containing the raw file data and metadata.
files = pw.io.gdrive.read(
object_id=os.getenv("FILE_OR_DIRECTORY_ID"),
service_user_credentials_file="secrets.json",
refresh_interval=30,
)
Parse, Chunk, and Vectorize Content
Incoming documents undergo parsing and embedding before storage. The implementation uses UnstructuredParser to extract text, TokenCountSplitter to create token-bounded chunks (40-120 tokens), and OpenAIEmbedder to generate 1536-dimensional vectors【L82-L86】. These vectors populate a KNNIndex that automatically refreshes when source files change, ensuring retrieval operates on the latest content.
Process Queries and Detect Alert Intent
Extract Alert Requests via LLM Classification
User queries arrive through an HTTP REST connector (pw.io.http.rest_connector) defined with QueryInputSchema. Before processing, the pipeline must determine if the user wants real-time alerts. The system calls an LLM with build_prompt_check_for_alert_request_and_extract_query to analyze the raw query text, returning a boolean flag indicating alert intent and the cleaned search query【L71-L79】.
query = (query
.select(prompt=build_prompt_check_for_alert_request_and_extract_query(pw.this.query))
.select(tupled=split_answer(OpenAIChat()(prompt_chat_single_qa(pw.this.prompt), max_tokens=100)))
.select(alert_enabled=pw.this.tupled[0], query=pw.this.tupled[1]))
Assign Persistent Query Identifiers
To correlate changing answers with original requests, each query receives a deterministic query_id generated by make_query_id based on the user and query content【L105-L108】. This identifier serves as the grouping key for stateful operations later in the pipeline.
Generate Context-Aware Answers
Retrieve Relevant Document Chunks
The system embeds the cleaned query and performs a nearest-neighbor search against the KNNIndex using index.get_nearest_items(query.data, k=3), retrieving the three most relevant document chunks【L38-L42】. These chunks provide the context window for answer generation.
Synthesize Responses with OpenAIChat
A second LLM call constructs the final answer using build_prompt, which concatenates the retrieved chunks with the user question【L44-L48】. The response flows through construct_message to wrap the answer text with the alert flag before returning to the API caller via response_writer【L58-L65】.
Trigger Real-Time Alerts on Content Changes
Compare Answer Versions with LLM Logic
When documents update, Pathway's reactive tables automatically propagate changes through the pipeline, re-evaluating answers for queries flagged with alert_enabled. An acceptor function compares the new answer against the previous version by calling build_prompt_compare_answers, which asks the LLM whether the responses deviate meaningfully, and interprets the "Yes/No" output via decision_to_bool【L73-L81】【L94-L101】.
def build_prompt_compare_answers(new: str, old: str) -> str:
return f"""
Are the two following responses deviating?
Answer with Yes or No.
First response: "{old}"
Second response: "{new}"
"""
Filter Duplicate Notifications
To prevent alert spam, pw.stateful.deduplicate filters the stream, retaining only distinct answer values per query_id【L82-L88】. This stateful operator maintains history across micro-batch iterations, ensuring users receive one notification per meaningful change rather than every minor revision.
Deliver Slack Notifications
Significant deviations trigger the final notification stage. The pipeline constructs a human-readable message via construct_notification_message, then dispatches it using pw.io.slack.send_alerts to the configured channel【L92-L95】. If Slack credentials are absent, the system falls back to console output for debugging.
alerts = (answer.filter(pw.this.alert_enabled)
.stateful.deduplicate(col=answer.response, acceptor=acceptor, instance=answer.query_id)
.select(message=construct_notification_message(pw.this.query, pw.this.response)))
pw.io.slack.send_alerts(alerts.message, os.getenv("SLACK_ALERT_CHANNEL_ID"), os.getenv("SLACK_ALERT_TOKEN"))
Complete Implementation Example
The following minimal snippet integrates all four stages—ingestion, query handling, answer generation, and alerting—into a single executable pipeline:
import pathway as pw
from pathway.xpacks.llm.embedders import OpenAIEmbedder
from pathway.xpacks.llm.llms import OpenAIChat, prompt_chat_single_qa
from pathway.xpacks.llm.parsers import UnstructuredParser
from pathway.xpacks.llm.splitters import TokenCountSplitter
from pathway.stdlib.ml.index import KNNIndex
import os
# 1️⃣ Ingest & index Google Drive files
files = pw.io.gdrive.read(
object_id=os.getenv("FILE_OR_DIRECTORY_ID"),
service_user_credentials_file="secrets.json",
refresh_interval=30,
)
docs = (files
.select(texts=UnstructuredParser()(pw.this.data))
.flatten(pw.this.texts)
.select(chunks=TokenCountSplitter()(pw.this.texts, min_tokens=40, max_tokens=120))
.flatten(pw.this.chunks))
enriched = docs + docs.select(data=OpenAIEmbedder()(pw.this.chunks))
index = KNNIndex(enriched.data, enriched, n_dimensions=1536)
# 2️⃣ Accept queries + detect alert intent
query, writer = pw.io.http.rest_connector(
host="0.0.0.0", port=8080, schema=QueryInputSchema,
autocommit_duration_ms=50, delete_completed_queries=False,
)
query = (query
.select(prompt=build_prompt_check_for_alert_request_and_extract_query(pw.this.query))
.select(tupled=split_answer(OpenAIChat()(prompt_chat_single_qa(pw.this.prompt), max_tokens=100)))
.select(user=pw.this.user, alert_enabled=pw.this.tupled[0], query=pw.this.tupled[1])
.select(data=OpenAIEmbedder()(pw.this.query),
query_id=pw.apply(make_query_id, pw.this.user, pw.this.query)))
# 3️⃣ Answer + send back
answer = (query
.join(index.get_nearest_items(query.data, k=3))
.select(prompt=build_prompt(pw.this.chunk, pw.this.query))
.select(response=OpenAIChat()(prompt_chat_single_qa(pw.this.prompt))))
writer(answer.select(result=construct_message(pw.this.response, pw.this.alert_enabled)))
# 4️⃣ Alert on changes
alerts = (answer.filter(pw.this.alert_enabled)
.stateful.deduplicate(col=answer.response, acceptor=acceptor, instance=answer.query_id)
.select(message=construct_notification_message(pw.this.query, pw.this.response)))
pw.io.slack.send_alerts(alerts.message, os.getenv("SLACK_CHANNEL_ID"), os.getenv("SLACK_TOKEN"))
pw.run()
Trigger an alert-enabled query using the REST API:
curl -X POST http://localhost:8080/ \
-H "Content-Type: application/json" \
-d '{"user":"alice","query":"When does the campaign start? Alert me if the date changes."}'
Summary
- Source Monitoring:
pw.io.gdrive.readpolls connected storage every 30 seconds to detect new or modified files intemplates/drive_alert/app.py. - Dynamic Indexing: KNNIndex automatically refreshes embeddings as documents change, ensuring retrievals always use current content.
- Alert Intent Detection: An initial LLM call via
build_prompt_check_for_alert_request_and_extract_querydetermines whether the user wants notifications for future changes. - Change Detection: The
acceptorfunction usesbuild_prompt_compare_answersto determine if new answers materially differ from previous versions before triggering notifications. - Deduplication:
pw.stateful.deduplicateprevents redundant alerts by tracking only distinct answer states perquery_id. - Delivery:
pw.io.slack.send_alertsdispatches notifications to Slack, with console fallback for local development.
Frequently Asked Questions
How does the system detect changes in Google Drive documents?
The pw.io.gdrive.read connector polls the Drive API every 30 seconds using the refresh_interval=30 parameter. When it detects file modifications, it emits updated rows into the Pathway table, triggering reactive recomputation of dependent queries and their associated alert logic.
What determines whether an alert is sent versus just an internal update?
Alerts require three conditions: the user must have enabled alerts via alert_enabled, the acceptor function must determine that the new answer meaningfully deviates from the previous version using build_prompt_compare_answers, and pw.stateful.deduplicate must confirm this specific answer state has not been sent before.
Can I implement real-time alerting with data sources other than Google Drive?
Yes. Replace pw.io.gdrive.read with any Pathway input connector such as pw.io.s3.read, pw.io.fs.read, or pw.io.sharepoint.read. The indexing, query handling, and alerting logic remains identical regardless of the source connector used.
How does the system prevent duplicate alerts for the same content change?
The pipeline uses pw.stateful.deduplicate with the query_id as the instance key. This operator maintains state across processing iterations, only emitting rows when the answer column changes to a value not previously seen for that specific query identifier.
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 →