Prerequisites
Before you begin, you will need:
pip install mixpeek
from mixpeek import Mixpeek client = Mixpeek(api_key="YOUR_API_KEY")
Step 1: Create a Namespace
A namespace is the top-level container in a multimodal data warehouse. It is analogous to a database in SQL or a collection in Qdrant. All collections, documents, and retrieval pipelines live within a namespace.
namespace = client.namespaces.create(
name="my-warehouse",
description="Production multimodal data warehouse",
)
print(f"Namespace created: {namespace['namespace_id']}")
# Every later call is scoped to this namespace through the client.
client = Mixpeek(api_key="YOUR_API_KEY", namespace=namespace["namespace_id"])Step 2: Define Collections
Collections are processing pipelines. Each collection defines how incoming objects are decomposed into features. You configure one feature extractor per collection.
Face Detection Collection
# Collections read from a bucket, so create the bucket they share first.
bucket = client.buckets.create(
bucket_name="raw-assets",
bucket_schema={"properties": {"video": {"type": "video"}}},
)
face_collection = client.collections.create(
collection_name="faces",
description="Face detection and recognition",
source={"type": "bucket", "bucket_ids": [bucket["bucket_id"]]},
feature_extractor={
"feature_extractor_name": "face_identity_extractor",
"version": "v1",
"parameters": {"detection_threshold": 0.8},
},
)Reference Logo Collection
# Mixpeek has no logo detector. Logo matching runs on image similarity: this
# collection embeds your reference logos with SigLIP, and a search compares
# frames against them.
logo_bucket = client.buckets.create(
bucket_name="reference-logos",
bucket_schema={"properties": {"logo": {"type": "image"}}},
)
logo_collection = client.collections.create(
collection_name="logos",
description="Reference logos for similarity matching",
source={"type": "bucket", "bucket_ids": [logo_bucket["bucket_id"]]},
feature_extractor={"feature_extractor_name": "image_extractor", "version": "v1"},
)Audio Fingerprint Collection
audio_collection = client.collections.create(
collection_name="audio",
description="Audio fingerprinting",
source={"type": "bucket", "bucket_ids": [bucket["bucket_id"]]},
feature_extractor={
"feature_extractor_name": "audio_fingerprint_extractor",
"version": "v1",
"parameters": {"segment_duration_sec": 5.0, "segment_hop_sec": 2.5},
},
)Text Extraction Collection
text_collection = client.collections.create(
collection_name="transcripts",
description="Speech-to-text transcription",
source={"type": "bucket", "bucket_ids": [bucket["bucket_id"]]},
feature_extractor={
"feature_extractor_name": "multimodal_extractor",
"version": "v1",
"parameters": {
"run_transcription": True,
"transcription_language": "en",
"run_transcription_embedding": True,
},
},
)Step 3: Ingest Objects
Ingestion in a multimodal data warehouse follows the pattern: upload to a bucket, then trigger the collections that read from it.
Upload Objects
# Land a video in the bucket the collections above read from. Each blob names
# the bucket_schema property it fills.
client.buckets.upload(
bucket["bucket_id"],
blobs=[{"property": "video", "type": "video", "data": "s3://raw-assets/video.mp4"}],
)Trigger Processing
A collection is bound to its source bucket when it is created. Processing starts when you trigger the collection, which batches the bucket objects it has not processed yet.
# Trigger each collection that reads the bucket. A trigger batches the objects
# that collection has not processed yet and returns the batch and its task.
runs = {name: client.collections.trigger(name) for name in ("faces", "audio", "transcripts")}Monitor Processing
After upload, objects are processed asynchronously. Monitor batch status:
# Check processing status through each batch's task
for name, run in runs.items():
task = client.tasks.get(run["task_id"])
print(f"{name}: batch {run['batch_id']} is {task['status']}")Step 4: Build Retrieval Pipelines
Multi-stage retrieval pipelines are the query layer of your warehouse. They compose filter, sort, reduce, enrich, and apply stages into expressive queries.
Basic Feature Search
retriever = client.retrievers.create(
retriever_name="face-search",
description="Search for faces by similarity",
collection_identifiers=["faces"],
input_schema={"face": {"type": "image", "required": True}},
stages=[
{
"stage_name": "match_faces",
"stage_id": "feature_search",
"parameters": {
"searches": [
{
"feature_uri": "mixpeek://face_identity_extractor@v1/insightface__arcface",
"query": {"input_mode": "content", "value": "{{INPUT.face}}"},
"top_k": 50,
}
],
"final_top_k": 50,
},
}
],
)
# Execute search with an image
results = client.retrievers.execute(
retriever["retriever_id"],
inputs={"face": "https://example.com/reference-face.jpg"},
)Multi-Stage Pipeline
pipeline = client.retrievers.create(
retriever_name="ip-safety-check",
description="Face matches, one per source video, with transcript context and IP risk labels",
collection_identifiers=["faces"],
input_schema={"face": {"type": "image", "required": True}},
stages=[
# Stage 1: Search for matching faces
{
"stage_name": "match_faces",
"stage_id": "feature_search",
"parameters": {
"searches": [
{
"feature_uri": "mixpeek://face_identity_extractor@v1/insightface__arcface",
"query": {"input_mode": "content", "value": "{{INPUT.face}}"},
"top_k": 100,
}
],
"final_top_k": 100,
},
},
# Stage 2: Put the scores on a common scale before anything reads them
{
"stage_name": "normalize",
"stage_id": "score_normalize",
"parameters": {"method": "min_max", "score_field": "score"},
},
# Stage 3: Keep one match per source video, then the top 20
{
"stage_name": "dedupe",
"stage_id": "deduplicate",
"parameters": {"strategy": "field", "fields": ["_internal.lineage.root_object_id"]},
},
{
"stage_name": "top_20",
"stage_id": "limit",
"parameters": {"limit": 20},
},
# Stage 4: Attach the transcript from the same source video
{
"stage_name": "add_transcript",
"stage_id": "document_enrich",
"parameters": {
"target_collection_id": text_collection["collection_id"],
"source_field": "_internal.lineage.root_object_id",
"target_field": "_internal.lineage.root_object_id",
"fields_to_merge": ["transcription"],
"output_field": "enrichments.transcript",
},
},
# Stage 5: Apply taxonomy classification
{
"stage_name": "classify_risk",
"stage_id": "taxonomy_enrich",
"parameters": {"taxonomy_id": "tax_your_ip_risk_taxonomy", "top_k": 1},
},
],
)Step 5: Apply Taxonomies
Taxonomies classify your unstructured data into structured categories. Configure them based on your use case.
Materialized Taxonomy (At Ingestion)
import requests
API = "https://api.mixpeek.com/v1"
HEADERS = {"Authorization": "Bearer YOUR_API_KEY", "X-Namespace": namespace["namespace_id"]}
# A flat taxonomy labels each document by matching it, through a retriever,
# against a collection of labeled examples. The SDK has no taxonomies
# resource, so this call is REST.
taxonomy = requests.post(API + "/taxonomies", headers=HEADERS, json={
"taxonomy_name": "content-type",
"description": "Classify content by type",
"config": {
"taxonomy_type": "flat",
"retriever_id": "ret_your_content_type_matcher",
"input_mappings": [{"input_key": "query", "source_type": "payload", "path": "transcription"}],
"source_collection": {"collection_id": "col_your_labeled_examples"},
},
}).json()
# Materialized: attach it to the collection, and documents are labeled as
# they are written.
client.collections.update(
"transcripts",
taxonomy_applications=[{"taxonomy_id": taxonomy["taxonomy_id"], "execution_mode": "materialize"}],
)On-Demand Taxonomy (At Query Time)
# On demand: run a taxonomy inside a retriever, so labels are computed for the
# results of each query and nothing is written back to the collection.
labeled_search = client.retrievers.create(
retriever_name="transcripts-with-sentiment",
collection_identifiers=["transcripts"],
input_schema={"query": {"type": "text", "required": True}},
stages=[
{
"stage_name": "search",
"stage_id": "feature_search",
"parameters": {
"searches": [
{
"feature_uri": "mixpeek://multimodal_extractor@v1/multilingual_e5_large_instruct_v1",
"query": {"input_mode": "text", "value": "{{INPUT.query}}"},
"top_k": 50,
}
],
"final_top_k": 50,
},
},
{
"stage_name": "brand_sentiment",
"stage_id": "taxonomy_enrich",
"parameters": {"taxonomy_id": "tax_your_brand_sentiment", "top_k": 1},
},
],
)Retroactive Taxonomy (Over Historical Data)
# Retroactive: apply a taxonomy that is already in the collection's
# taxonomy_applications to the documents written before it was attached.
requests.post(
API + "/collections/faces/apply-taxonomy",
headers=HEADERS,
json={"taxonomy_id": "tax_your_new_category_scheme"},
)Step 6: Configure Storage Tiering
Storage tiering moves a collection between the vector store and object storage.
# Move a collection out of the vector store. Its vectors stay searchable
# through object storage, and "active" brings them back.
requests.patch(
API + "/collections/faces/lifecycle",
headers=HEADERS,
json={"lifecycle_state": "cold"},
)Putting It Together: IP Safety Pipeline End-to-End
Here is a complete example that builds an IP safety pipeline from scratch:
from mixpeek import Mixpeek
client = Mixpeek(api_key="YOUR_API_KEY")
# 1. Create the namespace, then scope the client to it
ns = client.namespaces.create(name="ip-safety-prod")
client = Mixpeek(api_key="YOUR_API_KEY", namespace=ns["namespace_id"])
# 2. A bucket for the reference assets (the protected content to detect), and
# a collection per detection channel over it
bucket = client.buckets.create(
bucket_name="reference-assets",
bucket_schema={"properties": {"asset": {"type": "video"}}},
)
channels = {
"face-detection": {"feature_extractor_name": "face_identity_extractor", "version": "v1"},
"audio-fingerprint": {"feature_extractor_name": "audio_fingerprint_extractor", "version": "v1"},
}
for name, extractor in channels.items():
client.collections.create(
collection_name=name,
source={"type": "bucket", "bucket_ids": [bucket["bucket_id"]]},
feature_extractor=extractor,
)
# 3. Upload the reference assets, then process them
reference_asset_urls = ["s3://reference-assets/talent-a.mp4", "s3://reference-assets/jingle-b.mp4"]
for asset_url in reference_asset_urls:
client.buckets.upload(
bucket["bucket_id"],
blobs=[{"property": "asset", "type": "video", "data": asset_url}],
)
for name in channels:
client.collections.trigger(name)
# 4. The pre-publication check: faces in a frame of the new content, one match
# per reference asset, top 10
retriever = client.retrievers.create(
retriever_name="pre-pub-check",
collection_identifiers=["face-detection"],
input_schema={"frame": {"type": "image", "required": True}},
stages=[
{
"stage_name": "match_faces",
"stage_id": "feature_search",
"parameters": {
"searches": [
{
"feature_uri": "mixpeek://face_identity_extractor@v1/insightface__arcface",
"query": {"input_mode": "content", "value": "{{INPUT.frame}}"},
"top_k": 50,
}
],
"final_top_k": 50,
},
},
{"stage_name": "normalize", "stage_id": "score_normalize", "parameters": {"method": "min_max"}},
{"stage_name": "one_per_asset", "stage_id": "deduplicate", "parameters": {"strategy": "field", "fields": ["_internal.lineage.root_object_id"]}},
{"stage_name": "top_10", "stage_id": "limit", "parameters": {"limit": 10}},
],
)
# 5. Check new content before publication
results = client.retrievers.execute(
retriever["retriever_id"],
inputs={"frame": "https://example.com/new-content/frame-0120.jpg"},
)
matches = [doc for doc in results["documents"] if doc["score"] >= 0.8]
if matches:
print(f"IP conflicts detected: {len(matches)} matches")
for match in matches:
print(f" - {match['document_id']}: {match['score']:.2f}")
else:
print("Content cleared for publication")