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.read polls connected storage every 30 seconds to detect new or modified files in templates/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_query determines whether the user wants notifications for future changes.
  • Change Detection: The acceptor function uses build_prompt_compare_answers to determine if new answers materially differ from previous versions before triggering notifications.
  • Deduplication: pw.stateful.deduplicate prevents redundant alerts by tracking only distinct answer states per query_id.
  • Delivery: pw.io.slack.send_alerts dispatches 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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →