mirror of
https://github.com/qdrant/landing_page.git
synced 2026-10-09 21:08:31 +02:00
add incremental embedding updates tutorial
This commit is contained in:
@@ -6,5 +6,6 @@
|
||||
| [Time-Based Sharding](/documentation/tutorials-operations/time-based-sharding/) | Efficiently manage time-series data with user-defined sharding. | <span class="pill">Any</span> | 1h | <span class="text-yellow">Intermediate</span> |
|
||||
| [Large-Scale Search](/documentation/tutorials-operations/large-scale-search/) | Cost-efficient search for LAION-400M datasets. | <span class="pill">Any</span> | 48h | <span class="text-red">Advanced</span> |
|
||||
| [Secure a Self-Hosted Instance](/documentation/tutorials-operations/secure-qdrant/) | Enable TLS, API keys, and JWT access control. | <span class="pill">Any</span> | 45m | <span class="text-yellow">Intermediate</span> |
|
||||
| [Incremental Embedding Updates](/documentation/tutorials-operations/incremental-embedding-updates/) | Sync embeddings with changing raw text data. | <span class="pill">Python</span> | 25m | <span class="text-green">Beginner</span> |
|
||||
| [Qdrant Cloud Prometheus Monitoring](/documentation/ops-monitoring/managed-cloud-prometheus/) | Observability with Prometheus and Grafana. | <span class="pill">Prometheus</span> | 30m | <span class="text-yellow">Intermediate</span> |
|
||||
| [Self-Hosted Prometheus Monitoring](/documentation/ops-monitoring/hybrid-cloud-prometheus/) | Observability for hybrid/private cloud setups. | <span class="pill">Prometheus</span> | 30m | <span class="text-yellow">Intermediate</span> |
|
||||
+536
@@ -0,0 +1,536 @@
|
||||
---
|
||||
title: Incremental Embedding Updates
|
||||
short_description: "Sync embeddings with raw text data that changes over time."
|
||||
description: "Keep embeddings in the Qdrant search engine in sync with documentation that changes over time, for up-to-date vector search."
|
||||
weight: 32
|
||||
---
|
||||
|
||||
# Incremental Embedding Updates
|
||||
|
||||
| Time: 25 min | Level: Beginner | Output: [GitHub](https://github.com/qdrant/examples/blob/master/temporal-data-drift/sync_raw_data_to_embeddings.ipynb) | [](https://githubtocolab.com/qdrant/examples/blob/master/temporal-data-drift/sync_raw_data_to_embeddings.ipynb) |
|
||||
| --- | ----------- | ----------- | ----------- |
|
||||
|
||||
Qdrant documentation [lives on GitHub](https://github.com/qdrant/landing_page), consisting mainly of markdown pages with embedded code snippets and visuals.
|
||||
Like any other documentation of an evolving product, it's not static: raw data in markdowns changes with time, and users searching across our documentation expect to find the latest state of it.
|
||||
If search over documentation uses vectors, as ours does, it requires additional setup and maintenance to fulfill this expectation.
|
||||
|
||||
## Vectors <-> Raw Data
|
||||
|
||||
Vectors are a transformation of raw data.
|
||||
This transformation does not happen by itself when raw data changes. Unless vectors are updated proactively, documentation search would run against embeddings of text that no longer exists.
|
||||
There's a need for a re-embedding process, syncing vectors with raw data changes.
|
||||
|
||||
This tutorial provides a simple pipeline that, set up from day one, detects changes in your text data and executes incremental embedding updates.
|
||||
It reconciles a complete, current list of Qdrant documentation chunks with a Qdrant collection. Each run:
|
||||
|
||||
1. leaves unchanged chunks untouched,
|
||||
2. re-embeds changed text,
|
||||
3. reuses a vector when text changes location,
|
||||
4. adds new text,
|
||||
5. deletes text absent from the source list.
|
||||
|
||||
The pattern applies when your chunking is deterministic and enumerating the current source is inexpensive.
|
||||
|
||||
The tutorial has an accompanying [notebook](https://github.com/qdrant/examples/blob/master/temporal-data-drift/sync_raw_data_to_embeddings.ipynb).
|
||||
|
||||
## Prerequisites
|
||||
|
||||
```python
|
||||
%pip install -q "qdrant-client>=1.18"
|
||||
```
|
||||
|
||||
We use Qdrant Cloud and its [Free Embedding Inference](/documentation/cloud/inference/#free-embedding-models).
|
||||
Create a Free Tier [Qdrant Cloud cluster](https://cloud.qdrant.io/) and set `QDRANT_URL` and `QDRANT_API_KEY` in your environment.
|
||||
|
||||
|
||||
```python
|
||||
import os
|
||||
from qdrant_client import QdrantClient, models
|
||||
|
||||
QDRANT_URL = os.getenv("QDRANT_URL")
|
||||
QDRANT_API_KEY = os.getenv("QDRANT_API_KEY")
|
||||
|
||||
client = QdrantClient(
|
||||
url=QDRANT_URL,
|
||||
api_key=QDRANT_API_KEY,
|
||||
cloud_inference=True
|
||||
)
|
||||
```
|
||||
|
||||
## The Data: Qdrant Documentation
|
||||
|
||||
Let's look at the [operations tutorials](/documentation/tutorials-operations/) tab. Here's an example of a real change: this tutorial became a part of this tab, so our collection of vectors used for documentation search will have to be updated.
|
||||
|
||||
Let's consider a simple documentation hierarchy:
|
||||
|
||||
1. We have one page behind one `url`: https://qdrant.tech/documentation/tutorials-operations/secure-qdrant
|
||||
2. A page consists of sections. For example, the ["Step 2: Enable TLS" section](/documentation/tutorials-operations/secure-qdrant/#step-2-enable-tls). A section is marked by an `anchor`, generated from the heading text: "Step 2: Enable TLS" -> the "#step-2-enable-tls" part of the `section_url`.
|
||||
|
||||
Let's break down documentation using this hierarchy. Some sections might not fit the embedding model context window limit (how big of a text it can represent). We'll split them into chunks, numbered `0, 1, 2…`. For minimal hierarchy awareness, a chunk keeps its section heading prepended.
|
||||
|
||||
```text
|
||||
page https://qdrant.tech/documentation/tutorials-operations/secure-qdrant/ (url)
|
||||
├── section #prerequisites (anchor)
|
||||
│ └── chunk_num 0 "Prerequisites - Docker and Docker Compose installed..."
|
||||
├── section #secure-a-self-hosted-qdrant-instance (anchor)
|
||||
│ ├── chunk_num 0 "Secure a Self-Hosted Qdrant Instance | Time: 45 min..."
|
||||
│ └── chunk_num 1 "Secure a Self-Hosted Qdrant Instance > Qdrant Cloud..."
|
||||
├── section #step-1-start-an-unsecured-instance (anchor)
|
||||
│ └── chunk_num 0 "Step 1: Start an Unsecured Instance Start Qdrant..."
|
||||
├── section #step-2-enable-tls (anchor)
|
||||
│ └── chunk_num 0 "Step 2: Enable TLS Unencrypted connections allow..."
|
||||
└── ...
|
||||
```
|
||||
|
||||
So one page produces a set of chunks of the form: `{url, anchor, chunk_num, text}`. One vector = one section chunk.
|
||||
|
||||
<details>
|
||||
<summary>CHUNKS list used in this tutorial: three real tutorials from the operations tab chunked</summary>
|
||||
|
||||
```python
|
||||
CHUNKS = [ # three tutorials: secure-qdrant, migration, time-based-sharding
|
||||
{
|
||||
"url": "https://qdrant.tech/documentation/tutorials-operations/secure-qdrant/",
|
||||
"anchor": "prerequisites",
|
||||
"chunk_num": 0,
|
||||
"text": "Prerequisites - Docker and Docker Compose installed - `curl` available in your terminal - mkcert for generating a local self-signed certificate (installation instructions) - TLS requires Qdrant 1.2 or later, API key authentication requires Qdrant 1.2 or later, and granular access API keys (JWT) require Qdrant 1.9 or later. This tutorial uses the latest Qdrant image, which includes all these features. ---"
|
||||
},
|
||||
{
|
||||
"url": "https://qdrant.tech/documentation/tutorials-operations/secure-qdrant/",
|
||||
"anchor": "secure-a-self-hosted-qdrant-instance",
|
||||
"chunk_num": 0,
|
||||
"text": "Secure a Self-Hosted Qdrant Instance | Time: 45 min | Level: Intermediate | ..."
|
||||
},
|
||||
# ... full list in the ipynb
|
||||
]
|
||||
```
|
||||
|
||||
</details>
|
||||
|
||||
For each chunk we assume some text normalization pipeline is in place, as:
|
||||
|
||||
- Noise in the text degrades the embedding
|
||||
- Noise costs re-embedding when it's not needed (for example, someone added a trailing space)
|
||||
|
||||
```text
|
||||
normalize(text):
|
||||
- remove invisible characters (zero-width spaces, byte-order mark, soft hyphen)
|
||||
- collapse any whitespace run into a single space
|
||||
- ...
|
||||
```
|
||||
|
||||
## Configuring Collection
|
||||
|
||||
Let's configure a collection for chunks.
|
||||
|
||||
We'll use `sentence-transformers/all-MiniLM-L6-v2`: it's one of the [free embedding models](/documentation/cloud/inference/#free-embedding-models) on Qdrant Cloud Inference.
|
||||
Its output dimension is 384, its context window is 256 tokens, which is exactly why long sections got chunked above: over-window input is silently truncated.
|
||||
|
||||
### Collection Metadata
|
||||
|
||||
There are other types of drift harmful for production vector search, for example, a change in the embedding model version or in the data preparation pipeline.
|
||||
Vectors produced by different embedding models, or by the same model over differently prepared text, almost certainly should not mix in one collection: retrieval will degrade and it will be hard to detect why.
|
||||
|
||||
Let's consider a simple guardrail: save which model and which pipeline version produced the data points, in [**collection metadata**](/documentation/manage-data/collections/#collection-metadata), and verify against it. If one of the two changed, we need to trigger full collection re-embedding.
|
||||
|
||||
```python
|
||||
MODEL = "sentence-transformers/all-MiniLM-L6-v2"
|
||||
PIPELINE = "docs-prep-pipeline-v1"
|
||||
COLLECTION = "docs-sync-tutorial"
|
||||
|
||||
client.create_collection(
|
||||
COLLECTION,
|
||||
vectors_config=models.VectorParams(
|
||||
size=384, # all-MiniLM-L6-v2 output dimension
|
||||
distance=models.Distance.COSINE,
|
||||
),
|
||||
metadata={"embedding_model": MODEL, "pipeline_version": PIPELINE},
|
||||
)
|
||||
```
|
||||
|
||||
The gate against mixing embedding generations is then a simple check at the start of every run:
|
||||
|
||||
```text
|
||||
check_gate():
|
||||
read embedding_model and pipeline_version from collection metadata
|
||||
if either differs from this pipeline's constants:
|
||||
stop: full re-embed into a fresh collection required
|
||||
```
|
||||
|
||||
## Characteristics of a Document Chunk
|
||||
|
||||
What usually happens to documentation?
|
||||
Something completely new appears, information on pages gets fixed, pages get restructured and sections are moved as-is, pages get deleted.
|
||||
|
||||
It makes sense to monitor two independent characteristics of a document chunk:
|
||||
|
||||
- **Content**: the text we search against and generate the embedding from.
|
||||
- **Position**: where the chunk lives, in our case its URL, anchor, and number.
|
||||
|
||||
Hence every record should get two derived values:
|
||||
|
||||
- **Content fingerprint**, like SHA-256 of the text. It changes if a single character changes, and never otherwise. Comparing fingerprints answers "*Is it the same content?*" without comparing texts.
|
||||
- **Deterministic ID** for position in documentation. For example, `url + "#" + anchor + "::" + chunk_num` turned into a UUID, one of the two point ID formats Qdrant accepts. Comparing IDs answers "*Is this content still at the same position?*".
|
||||
|
||||
```python
|
||||
import hashlib
|
||||
import uuid
|
||||
from datetime import datetime, timezone
|
||||
|
||||
def content_hash(text):
|
||||
return hashlib.sha256(text.encode()).hexdigest()
|
||||
|
||||
def point_id(url, anchor, num):
|
||||
# NAMESPACE_URL is a fixed constant uuid5 requires; it marks the input as a URL-like name
|
||||
return str(uuid.uuid5(uuid.NAMESPACE_URL, f"{url}#{anchor}::{num}"))
|
||||
|
||||
def prepare_chunks_for_sync(chunks):
|
||||
"""Derive both values (and the section address) for every raw chunk."""
|
||||
out = []
|
||||
for c in chunks:
|
||||
text = normalize(c["text"])
|
||||
out.append({
|
||||
**c,
|
||||
"text": text,
|
||||
"section_url": f"{c['url']}#{c['anchor']}" if c["anchor"] else c["url"],
|
||||
"content_hash": content_hash(text),
|
||||
"point_id": point_id(c["url"], c["anchor"], c["chunk_num"]),
|
||||
})
|
||||
return out
|
||||
```
|
||||
Example:
|
||||
|
||||
```text
|
||||
point ID: 2ff5204a-0353-5991-... # UUID(url + "#" + anchor + "::" + chunk_num)
|
||||
text (to vectorize): Prerequisites - Docker and Docker Compose...
|
||||
content_hash: 27d55e75b962f1d5... # sha256(text)
|
||||
```
|
||||
|
||||
Additionally, a point can be described by the following fields:
|
||||
|
||||
**Payload:**
|
||||
- `url`: filter or group all chunks of one page
|
||||
- `section_url`: filter or group all chunks of one section
|
||||
- `last_updated`: when content of this chunk last changed (or was created)
|
||||
|
||||
<details>
|
||||
<summary>payload() implementation</summary>
|
||||
|
||||
```python
|
||||
def payload(chunk, last_updated=None):
|
||||
return {
|
||||
"url": chunk["url"],
|
||||
"anchor": chunk["anchor"],
|
||||
"chunk_num": chunk["chunk_num"],
|
||||
"section_url": chunk["section_url"],
|
||||
"text": chunk["text"],
|
||||
"content_hash": chunk["content_hash"],
|
||||
"last_updated": last_updated or datetime.now(timezone.utc).isoformat(timespec="seconds"),
|
||||
}
|
||||
```
|
||||
|
||||
</details>
|
||||
|
||||
For all the payload fields used for filtering or grouping we need to create a [**payload index**](/documentation/manage-data/indexing/).
|
||||
|
||||
```python
|
||||
for field in ("content_hash", "url", "section_url"):
|
||||
client.create_payload_index(COLLECTION, field, models.PayloadSchemaType.KEYWORD)
|
||||
```
|
||||
|
||||
## Populate Collection
|
||||
|
||||
Populate the collection with the whole documentation.
|
||||
|
||||
```python
|
||||
client.upsert(COLLECTION, points=[
|
||||
models.PointStruct(
|
||||
id=c["point_id"],
|
||||
vector=models.Document(text=c["text"], model=MODEL), # Cloud Inference embeds text server-side
|
||||
payload=payload(c),
|
||||
)
|
||||
for c in prepare_chunks_for_sync(CHUNKS)
|
||||
], wait=True)
|
||||
```
|
||||
|
||||
<details>
|
||||
<summary>Test the search against it</summary>
|
||||
|
||||
```python
|
||||
QUERY = "Where exactly to set `QDRANT__SERVICE__API_KEY` variable to enable authentication for a self-hosted Qdrant?"
|
||||
|
||||
client.query_points(
|
||||
COLLECTION,
|
||||
query=models.Document(text=QUERY, model=MODEL),
|
||||
limit=3,
|
||||
with_payload=["section_url", "text"],
|
||||
)
|
||||
```
|
||||
|
||||
You should get something like:
|
||||
|
||||
```text
|
||||
0.675 https://qdrant.tech/documentation/tutorials-operations/secure-qdrant/#secure-a-self-hosted-qdrant-instance
|
||||
Secure a Self-Hosted Qdrant Instance | Time: 45 min | Level: Intermediate | ...
|
||||
```
|
||||
|
||||
</details>
|
||||
|
||||
## Syncing with Documentation Changes
|
||||
|
||||
Your sync trigger could be a CI job on merge if your docs live in git, or a nightly cron job.
|
||||
|
||||
The input of every sync with a documentation collection here is the **current full chunk list of the docs**. For a simple deterministic data prep pipeline it's cheap to gather this full list once a day, saving you the headache of deriving raw changes.
|
||||
|
||||
Each incoming chunk is compared against the current documentation collection in one of the following ways, based on the `point ID` (the chunk's address in documentation) and `content_hash` (the chunk's exact content, its fingerprint):
|
||||
|
||||
```text
|
||||
incoming chunk
|
||||
├─ ID found in the collection?
|
||||
│ ├─ yes: fingerprint equal?
|
||||
│ │ ├─ yes -> unchanged: the point stays as is
|
||||
│ │ └─ no -> content changed: re-embed in place
|
||||
│ └─ no: identical fingerprint under another ID?
|
||||
│ ├─ yes -> address changed: reuse the vector, create a new point
|
||||
│ └─ no -> new: embed and insert a new point
|
||||
└─ stored point whose ID is absent from the incoming list
|
||||
-> gone: delete (last, after all writes)
|
||||
```
|
||||
|
||||
**Note:** Optional footgun-guard: take a [snapshot](/documentation/snapshots/) before sync, delete it later when everything looks fine.
|
||||
|
||||
### Input of a Sync Pipeline
|
||||
|
||||
Let's consider some possible changes:
|
||||
|
||||
- adding to the ["Secure a Self-Hosted Qdrant Instance" tutorial](/documentation/tutorials-operations/secure-qdrant/) a new small section "Step 6: Rotate API keys",
|
||||
with the "Step 3" section now pointing to it
|
||||
- the [migration page](/documentation/tutorials-operations/migration/) moved to a new URL
|
||||
- the ["Time-based sharding" tutorial](/documentation/tutorials-operations/time-based-sharding/) was removed
|
||||
|
||||
<details>
|
||||
<summary>The LATEST_CHUNKS list with these three changes</summary>
|
||||
|
||||
```python
|
||||
untouched_secure_qdrant = [
|
||||
c for c in CHUNKS
|
||||
if c["url"] == "https://qdrant.tech/documentation/tutorials-operations/secure-qdrant/"
|
||||
and c["anchor"] != "step-3-enable-an-admin-api-key"
|
||||
]
|
||||
|
||||
# now points to the new section
|
||||
step_3 = {
|
||||
"url": "https://qdrant.tech/documentation/tutorials-operations/secure-qdrant/",
|
||||
"anchor": "step-3-enable-an-admin-api-key",
|
||||
"chunk_num": 0,
|
||||
"text": "Step 3: Enable an Admin API Key ... Refer to Security > Authentication to learn more about admin API keys, including API key rotation. --- See also: rotating API keys."
|
||||
}
|
||||
|
||||
# the new section
|
||||
step_6 = {
|
||||
"url": "https://qdrant.tech/documentation/tutorials-operations/secure-qdrant/",
|
||||
"anchor": "step-6-rotate-api-keys",
|
||||
"chunk_num": 0,
|
||||
"text": "Step 6: Rotate API keys Rotate the admin API key on a schedule and immediately after any suspected exposure. Update every client before revoking the old key."
|
||||
}
|
||||
|
||||
# the migration page moved: same texts, new addresses
|
||||
moved = [
|
||||
{**c, "url": "https://qdrant.tech/documentation/tutorials-operations/migration-guide/"}
|
||||
for c in CHUNKS
|
||||
if c["url"] == "https://qdrant.tech/documentation/tutorials-operations/migration/"
|
||||
]
|
||||
|
||||
# the time-based-sharding tutorial is absent from LATEST_CHUNKS - that is how a deletion arrives
|
||||
|
||||
LATEST_CHUNKS = prepare_chunks_for_sync(untouched_secure_qdrant + [step_3, step_6] + moved)
|
||||
```
|
||||
|
||||
</details>
|
||||
|
||||
We now check every incoming chunk against the collection: does its ID (address) exist, and does its `content_hash` (exact text) match?
|
||||
|
||||
[`retrieve`](/documentation/manage-data/points/) fetches points by ID. At corpus scale you would batch the IDs.
|
||||
|
||||
```python
|
||||
def split_by_state(latest_chunks):
|
||||
"""Compare the incoming chunk list to the collection: who is unchanged, changed, or unknown."""
|
||||
incoming = {c["point_id"]: c for c in latest_chunks}
|
||||
|
||||
stored = {}
|
||||
points = client.retrieve(
|
||||
COLLECTION,
|
||||
ids=list(incoming),
|
||||
with_payload=["content_hash"],
|
||||
with_vectors=False,
|
||||
)
|
||||
for p in points:
|
||||
stored[str(p.id)] = p.payload["content_hash"]
|
||||
|
||||
unchanged, content_changed, unknown_ids = [], [], []
|
||||
for pid, c in incoming.items():
|
||||
if stored.get(pid) == c["content_hash"]:
|
||||
unchanged.append(c)
|
||||
elif pid in stored:
|
||||
content_changed.append(c)
|
||||
else:
|
||||
unknown_ids.append(c)
|
||||
|
||||
return incoming, unchanged, content_changed, unknown_ids
|
||||
|
||||
|
||||
incoming_ids, unchanged, content_changed, unknown_ids = split_by_state(LATEST_CHUNKS)
|
||||
```
|
||||
|
||||
### Case 1: Unchanged, Do Nothing
|
||||
|
||||
These chunks carry the same fingerprint as before.
|
||||
|
||||
### Case 2: Content Changed, Re-Embed
|
||||
|
||||
The chunk about Step 3 exists under a known ID (it didn't change its position on the docs website) but carries new information.
|
||||
Use `upsert`: writing a point under an existing ID replaces it.
|
||||
|
||||
```python
|
||||
def re_embed_changed(content_changed):
|
||||
if not content_changed:
|
||||
return
|
||||
client.upsert(COLLECTION,
|
||||
points=[
|
||||
models.PointStruct(
|
||||
id=c["point_id"],
|
||||
vector=models.Document(text=c["text"], model=MODEL),
|
||||
payload=payload(c),
|
||||
)
|
||||
for c in content_changed],
|
||||
wait=True)
|
||||
```
|
||||
|
||||
### Cases 3 and 4: ID Is Not Present in the Collection
|
||||
|
||||
Six IDs are unknown to the collection, but an unknown ID does not necessarily mean new content. When a page moves as-is, every chunk on it gets a new address (a new ID), while the text stays exactly the same. Embedding it again would produce the same vector, so why pay for it.
|
||||
|
||||
A filtered [`scroll`](/documentation/manage-data/points/) on `content_hash` answers the question "*does this exact text already exist under some other ID?*".
|
||||
- On a hit, we copy the stored vector into the new point and keep the source's `last_updated` as the content did not change.
|
||||
- On a miss, the content is genuinely new; we embed and insert a new point.
|
||||
|
||||
**Note:** *This version performs one hash lookup per unknown chunk so the decision is easy to inspect. In production, batch hash lookups and point upserts.*
|
||||
|
||||
```python
|
||||
def reuse_or_add(unknown_ids):
|
||||
"""Reuse an existing embedding when the same text is already stored; embed only what is new."""
|
||||
reused, added = 0, 0
|
||||
|
||||
for c in unknown_ids:
|
||||
same_text = models.Filter(must=[
|
||||
models.FieldCondition(
|
||||
key="content_hash",
|
||||
match=models.MatchValue(value=c["content_hash"]),
|
||||
)
|
||||
])
|
||||
hits, _ = client.scroll(
|
||||
COLLECTION,
|
||||
scroll_filter=same_text,
|
||||
limit=1,
|
||||
with_payload=["last_updated"],
|
||||
with_vectors=True,
|
||||
)
|
||||
|
||||
if hits: # same text, new address: copy the vector, keep its last_updated
|
||||
point = models.PointStruct(
|
||||
id=c["point_id"],
|
||||
vector=hits[0].vector,
|
||||
payload=payload(c, hits[0].payload["last_updated"]),
|
||||
)
|
||||
reused += 1
|
||||
else: # genuinely new content: embed and insert
|
||||
point = models.PointStruct(
|
||||
id=c["point_id"],
|
||||
vector=models.Document(text=c["text"], model=MODEL),
|
||||
payload=payload(c),
|
||||
)
|
||||
added += 1
|
||||
|
||||
client.upsert(COLLECTION, points=[point], wait=True)
|
||||
|
||||
return reused, added
|
||||
```
|
||||
|
||||
What's important to notice: the old points, the migration page under its old URL, are still in the collection. They need to be removed, and that is the last case.
|
||||
|
||||
### Case 5: Gone, Delete
|
||||
|
||||
Whatever LATEST_CHUNKS does not contain no longer exists at the source. The deletion is one filtered call, "every point whose ID is *not* in the incoming list".
|
||||
|
||||
**Note:** *Deletion runs **last**, after all writes: hence if someone queries documentation at night while this pipeline runs, results for a second might be weird, as mid-run search here sees old and new content side by side:)*
|
||||
|
||||
**Note:** *It's a good practice to put some guardrails on the number of deletions before running it: if it is suspiciously large, you might want to skip deletion and investigate instead. Mind the edge case: an empty incoming list would match every point in the collection, so refuse to sync empty input.*
|
||||
|
||||
```python
|
||||
def delete_gone(incoming_ids):
|
||||
"""Remove every point the current crawl no longer contains. Returns how many."""
|
||||
if not incoming_ids:
|
||||
raise ValueError("Refusing to delete from an empty source snapshot.")
|
||||
|
||||
stale = models.Filter(must_not=[models.HasIdCondition(has_id=list(incoming_ids))])
|
||||
|
||||
to_delete = client.count(COLLECTION, count_filter=stale).count
|
||||
|
||||
# potential check against a threshold to avoid accidental mass deletion could be added here
|
||||
client.delete(COLLECTION, points_selector=models.FilterSelector(filter=stale), wait=True)
|
||||
return to_delete
|
||||
```
|
||||
|
||||
## Run and Verify the Sync
|
||||
|
||||
The five cases, assembled from the functions defined above:
|
||||
|
||||
```python
|
||||
def sync(latest_chunks):
|
||||
check_gate() # refuse to mix embedding models or pipeline versions
|
||||
|
||||
chunks = prepare_chunks_for_sync(latest_chunks)
|
||||
incoming_ids, unchanged, content_changed, unknown_ids = split_by_state(chunks)
|
||||
|
||||
re_embed_changed(content_changed)
|
||||
reused, added = reuse_or_add(unknown_ids)
|
||||
deleted = delete_gone(incoming_ids)
|
||||
|
||||
return {
|
||||
"unchanged": len(unchanged),
|
||||
"re-embedded": len(content_changed),
|
||||
"reused_embedding": reused,
|
||||
"added": added,
|
||||
"deleted": deleted,
|
||||
}
|
||||
```
|
||||
|
||||
Run the sync.
|
||||
|
||||
```python
|
||||
run = sync(LATEST_CHUNKS)
|
||||
print(run)
|
||||
```
|
||||
|
||||
You should see something like:
|
||||
|
||||
```text
|
||||
{'unchanged': 9, 're-embedded': 1, 'reused_embedding': 5, 'added': 1, 'deleted': 20}
|
||||
```
|
||||
|
||||
A re-run of the same sync input should change nothing: every change counter at zero, all 16 chunks in `unchanged`.
|
||||
|
||||
## Conclusion
|
||||
|
||||
A deterministic ID, a content fingerprint, and five sync cases keep embeddings in sync with changing raw data, re-embedding only what actually changed. Adapt this pipeline to your own documents.
|
||||
|
||||
Ways to make it better:
|
||||
|
||||
- Pipelines that risk concurrent iterative updates: see [conditional updates](/documentation/manage-data/points/#conditional-updates) and [update modes](/documentation/manage-data/points/#update-mode), per-write preconditions. This pipeline runs one sync at a time and does not need them.
|
||||
- The `last_updated` field this sync maintains can power recency-aware ranking via [decay functions in a formula query](/documentation/search/search-relevance/).
|
||||
|
||||
Related guides:
|
||||
|
||||
- Switching or upgrading the embedding model: [Embedding Model Migration](/documentation/tutorials-operations/embedding-model-migration/)
|
||||
- Wholesale infrastructure swaps: [Blue-Green Deployment](https://qdrant.tech/documentation/tutorials-operations/blue-green-deployment/)
|
||||
- Sync driven by database change events: [Data Synchronization](/documentation/data-synchronization/)
|
||||
Reference in New Issue
Block a user