# How to Implement Real-Time Alerting When Document Content Changes in Pathway

> Implement real-time alerting for document changes with Pathway. Get instant Slack notifications when query answers change in connected sources like Google Drive using the drive_alert template.

- Repository: [Pathway/llm-app](https://github.com/pathwaycom/llm-app)
- Tags: how-to-guide
- Published: 2026-03-07

---

**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`](https://github.com/pathwaycom/llm-app/blob/main/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.

```python
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】.

```python
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】.

```python
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.

```python
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:

```python
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:

```bash
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`](https://github.com/pathwaycom/llm-app/blob/main/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.