Scaffold locally · batch
Media indexing
Batch ingest from object storage, process every modality, export to your warehouse.
Scaffold locally
uvx pixeltable-new --template media-indexing my-media-indexingSame starter-kit files Cloud uses. Local UI (static HTML in some templates) is for uvx, not for Cloud. Cloud deploys schema + insert routes via pxt serve.
schema.py
"""Enterprise media processing pipeline -- ingest, process, search across modalities."""
import pixeltable as pxt
from pixeltable.functions import image as pxt_image
from pixeltable.functions.audio import audio_splitter
from pixeltable.functions.document import document_splitter
from pixeltable.functions.huggingface import sentence_transformer
from pixeltable.functions.uuid import uuid7
NAMESPACE = "pipeline"
EMBED_MODEL = "all-MiniLM-L6-v2"
embed_fn = sentence_transformer.using(model_id=EMBED_MODEL)
pxt.create_dir(NAMESPACE, if_exists="ignore")
# ---------------------------------------------------------------------------
# Media ingest table -- central registry of all incoming media
# ---------------------------------------------------------------------------
media = pxt.create_table(
f"{NAMESPACE}.media",
{
"url": pxt.String,
"media_type": pxt.String, # 'image' | 'document' | 'audio' | 'video'
"tags": pxt.Json,
"uuid": uuid7(),
"timestamp": pxt.Timestamp,
},
primary_key=["uuid"],
if_exists="ignore",
)
# ---------------------------------------------------------------------------
# Image processing
# ---------------------------------------------------------------------------
images = pxt.create_table(
f"{NAMESPACE}.images",
{
"image": pxt.Image,
"source_url": pxt.String,
"uuid": uuid7(),
"timestamp": pxt.Timestamp,
},
primary_key=["uuid"],
if_exists="ignore",
)
images.add_computed_column(
thumbnail=pxt_image.b64_encode(pxt_image.thumbnail(images.image, size=(320, 320))),
if_exists="ignore",
)
images.add_computed_column(width=images.image.width, if_exists="ignore")
images.add_computed_column(height=images.image.height, if_exists="ignore")
images.add_computed_column(mode=images.image.mode, if_exists="ignore")
# Optional: CLIP embedding for visual search (uncomment to enable)
# from pixeltable.functions.huggingface import clip
# clip_fn = clip.using(model_id='openai/clip-vit-base-patch32')
# images.add_embedding_index('image', embedding=clip_fn, if_exists='ignore')
# Optional: vision LLM caption (uncomment + set OPENAI_API_KEY)
# from pixeltable.functions.openai import chat_completions
# images.add_computed_column(
# caption=chat_completions(
# messages=[{'role': 'user', 'content': [
# {'type': 'text', 'text': 'Describe this image in one sentence.'},
# {'type': 'image_url', 'image_url': {'url': images.image}},
# ]}],
# model='gpt-4o-mini',
# ).choices[0].message.content,
# if_exists='ignore',
# )
# Cloud storage: persist thumbnails to S3
# images.add_computed_column(thumbnail=..., media_destination='s3://bucket/thumbnails/')
# ---------------------------------------------------------------------------
# Document processing
# ---------------------------------------------------------------------------
documents = pxt.create_table(
f"{NAMESPACE}.documents",
{
"document": pxt.Document,
"title": pxt.String,
"uuid": uuid7(),
"timestamp": pxt.Timestamp,
},
primary_key=["uuid"],
if_exists="ignore",
)
doc_chunks = pxt.create_view(
f"{NAMESPACE}.doc_chunks",
documents,
iterator=document_splitter(
documents.document,
separators="sentence,token_limit",
limit=300,
metadata="page",
),
if_exists="ignore",
)
doc_chunks.add_embedding_index("text", idx_name="doc_text_idx", string_embed=embed_fn, if_exists="ignore")
# Optional: LLM summary per document (uncomment + set OPENAI_API_KEY)
# documents.add_computed_column(
# summary=chat_completions(
# messages=[{'role': 'user', 'content': 'Summarize:\n' + documents.title}],
# model='gpt-4o-mini',
# ).choices[0].message.content,
# if_exists='ignore',
# )
# ---------------------------------------------------------------------------
# Audio processing
# ---------------------------------------------------------------------------
audio_files = pxt.create_table(
f"{NAMESPACE}.audio_files",
{
"audio": pxt.Audio,
"title": pxt.String,
"uuid": uuid7(),
"timestamp": pxt.Timestamp,
},
primary_key=["uuid"],
if_exists="ignore",
)
audio_chunks = pxt.create_view(
f"{NAMESPACE}.audio_chunks",
audio_files,
iterator=audio_splitter(audio_files.audio, duration=30.0),
if_exists="ignore",
)
# Optional: transcription (uncomment + install whisper or set OPENAI_API_KEY)
# from pixeltable.functions.openai import transcriptions
# audio_chunks.add_computed_column(
# transcript=transcriptions(audio=audio_chunks.audio_segment, model='whisper-1'),
# if_exists='ignore',
# )
# Version control: snapshot before destructive changes
# pxt.create_snapshot('pipeline.snapshot_v1', media, if_exists='ignore')
# ---------------------------------------------------------------------------
# Query functions -- used by pxt serve and pipeline.py
# ---------------------------------------------------------------------------
@pxt.query
def search_documents(query_text: str, limit: int = 10):
"""Semantic search over document chunks."""
sim = doc_chunks.text.similarity(string=query_text)
return (
doc_chunks.order_by(sim, asc=False).limit(limit).select(text=doc_chunks.text, score=sim, page=doc_chunks.page)
)
@pxt.query
def list_images():
"""List all processed images with metadata."""
return images.select(
uuid=images.uuid,
source_url=images.source_url,
width=images.width,
height=images.height,
mode=images.mode,
timestamp=images.timestamp,
).order_by(images.timestamp, asc=False)
@pxt.query
def list_documents():
"""List all documents."""
return documents.select(
uuid=documents.uuid,
title=documents.title,
document=documents.document,
timestamp=documents.timestamp,
).order_by(documents.timestamp, asc=False)
if __name__ == "__main__":
print("Schema initialized. Run: pxt serve pipeline")
README
Content Pipeline -- Enterprise Media Processing
Ingest media from S3/URLs, auto-process across modalities, export structured results to your database. Your own Cloudinary AI processing layer, self-hosted.
What it replaces: Cloudinary AI Transform ($5K--50K/yr), Azure Content Understanding pipelines, custom media processing microservices.
What You Get
| Modality | Processing | Search |
|---|---|---|
| Images | Thumbnails, dimensions, mode, optional CLIP embedding + vision caption | Visual similarity (CLIP) |
| Documents | Sentence + token chunking, page metadata, embedding index | Semantic search over chunks |
| Audio | 30s segment splitting, optional Whisper transcription | -- |
Quickstart
uv sync # install deps
uv run python schema.py # initialize tables
uv run pxt serve pipeline # http://localhost:8000/docs
Or run batch processing instead:
uv run python pipeline.py --urls path/to/image.png path/to/report.pdf path/to/audio.mp3
uv run python pipeline.py --status
Two Modes
Real-time API (pxt serve)
Routes are declared in pyproject.toml -- zero web code required:
POST /api/search -- semantic search over document chunks
POST /api/ingest/image -- upload + process an image
POST /api/ingest/document -- upload + process a document
POST /api/ingest/audio -- upload + process audio
Batch Processing (python pipeline.py)
python pipeline.py --file urls.json # ingest from JSON
python pipeline.py --urls s3://bucket/img.png # ingest from CLI
python pipeline.py --search "quarterly revenue" # search documents
python pipeline.py --export-parquet output/ # export to Parquet
python pipeline.py --status # show counts
Cloud I/O Patterns
Importing from S3
Pixeltable resolves S3 URLs natively -- just pass them as source URLs:
images.insert([{'image': 's3://my-bucket/photos/hero.jpg', 'source_url': 's3://my-bucket/photos/hero.jpg', 'timestamp': datetime.now()}])
documents.insert([{'document': 's3://my-bucket/docs/report.pdf', 'title': 'Q4 Report', 'timestamp': datetime.now()}])
Set AWS_ACCESS_KEY_ID and AWS_SECRET_ACCESS_KEY in your environment (or use IAM roles).
Exporting to Postgres / Snowflake
from pixeltable.io.sql import export_sql
export_sql(
images.select(images.uuid, images.source_url, images.width, images.height),
'processed_images',
db_connect_str='postgresql://user:pass@host:5432/mydb',
if_exists='replace',
)
export_sql(
doc_chunks.select(doc_chunks.text, doc_chunks.page),
'document_chunks',
db_connect_str='snowflake://user:pass@account/db/schema',
if_exists='replace',
)
Exporting to Parquet
python pipeline.py --export-parquet output/
Project Structure
media-indexing/
├── schema.py -- tables, views, computed columns, query functions
├── pipeline.py -- batch runner (CLI alternative to pxt serve)
├── pyproject.toml -- dependencies + pxt serve route config
└── README.md
Optional Features
Uncomment sections in schema.py to enable:
- Vision captions -- GPT-4o-mini image descriptions (requires
OPENAI_API_KEY) - CLIP visual search -- image similarity search via
openai/clip-vit-base-patch32 - Audio transcription -- Whisper transcription on audio segments (requires
OPENAI_API_KEY) - Document summaries -- LLM-generated summaries per document
- Cloud storage -- persist computed media to S3 via
destination='s3://...'
Version Control
Snapshot your pipeline state before destructive changes:
import pixeltable as pxt
pxt.create_snapshot('pipeline.snapshot_v1', media, if_exists='ignore')