---
title: "Beyond AVG(): Building Custom Aggregation Functions for AI Workflows with Pixeltable UDA"
date: "2025-10-12"
author: "Pixeltable Team"
tags:
  - Custom Aggregations
  - UDA
  - User-Defined Aggregates
  - AI Metrics
  - Video Analytics
  - Embedding Analytics
  - Advanced Pixeltable
  - Data Aggregation
  - Python UDFs
description: "SQL's AVG() and COUNT() don't work for embeddings, video frames, or multimodal data. Learn how to build custom aggregation functions with Pixeltable's @pxt.uda decorator for specialized AI metrics, rolling averages over embeddings, and domain-specific analytics."
url: "https://pixeltable.com/blog/beyond-avg-custom-aggregations-uda"
---

# Beyond AVG(): Building Custom Aggregation Functions for AI Workflows with Pixeltable UDA

## The Aggregation Problem: When SQL Functions Don't Cut It

 
You're building an AI application that processes thousands of video frames. You need to calculate specialized metrics: average embedding similarity across frames, detect scene changes based on visual variance, or compute rolling statistics over multimodal data. You reach for SQL's standard aggregation functions (`AVG()`, `COUNT()`, `SUM()`) and hit a wall.

 
 
**The problem:** Traditional SQL aggregations weren't designed for AI data types. They work great for numbers and simple statistics, but try to aggregate over embeddings, calculate custom video metrics, or build domain-specific analytics, and you're forced into awkward workarounds or manual post-processing.

 
 
This is where **User-Defined Aggregates (UDA)** transform AI development. With Pixeltable's `@pxt.uda` decorator, you can build custom aggregation functions that understand your AI data and operate efficiently at scale.

 
## What Are User-Defined Aggregates?

 
A User-Defined Aggregate (UDA) is a custom function that combines multiple rows of data into a single result, just like SQL's built-in `SUM()` or `AVG()`, but with logic you define for your specific AI use case.

 
 
Unlike simple [User-Defined Functions (UDFs)](/blog/python-udfs-pixeltable) that operate on single rows, UDAs accumulate state across multiple rows and produce a final aggregated result. This is perfect for AI workflows where you need custom metrics over groups of images, frames, embeddings, or any multimodal data.

 
### Anatomy of a Pixeltable UDA

 
Every UDA implements three key methods:

 

 - **`__init__()`:** Initialize the aggregator with starting state

 - **`update()`:** Process each row and update internal state

 - **`value()`:** Return the final aggregated result

 

 
```python

import pixeltable as pxt
from pixeltable import Aggregator

@pxt.uda
class CustomAverage(Aggregator):
 """Simple custom average aggregation"""
 
 def __init__(self):
 """Initialize aggregator state"""
 self.sum = 0.0
 self.count = 0
 
 def update(self, value: float):
 """Update state with each row"""
 if value is not None:
 self.sum += value
 self.count += 1
 
 def value(self) -> float:
 """Return final aggregated result"""
 return self.sum / self.count if self.count > 0 else 0.0

# Use like any built-in aggregation
result = table.select(CustomAverage(table.score)).collect()
 
```

 
## Real-World AI Aggregations: Where UDAs Shine

 
### 1. Average Embedding Similarity Across Frames

 
Calculate how visually similar video frames are within scenes:

 
```python

import pixeltable as pxt
from pixeltable import Aggregator
import numpy as np

@pxt.uda
class AverageEmbeddingSimilarity(Aggregator):
 """Calculate average cosine similarity across embeddings"""
 
 def __init__(self):
 self.embeddings = []
 
 def update(self, embedding: np.ndarray):
 """Accumulate embeddings"""
 if embedding is not None and len(embedding) > 0:
 self.embeddings.append(embedding)
 
 def value(self) -> float:
 """Compute average pairwise similarity"""
 if len(self.embeddings) dict:
 """Return variance statistics for scene analysis"""
 if not self.variance_scores:
 return {'mean_variance': 0.0, 'max_variance': 0.0, 'scene_changes': 0}
 
 variance_array = np.array(self.variance_scores)
 mean_var = float(np.mean(variance_array))
 max_var = float(np.max(variance_array))
 
 # Detect scene changes (large jumps in variance)
 threshold = mean_var + 2 * np.std(variance_array)
 scene_changes = int(np.sum(variance_array > threshold))
 
 return {
 'mean_variance': mean_var,
 'max_variance': max_var,
 'scene_changes': scene_changes,
 'variance_scores': self.variance_scores
 }

# Analyze video scene changes
scene_analysis = frames.select(
 frames.video_id,
 analysis=VisualVarianceDetector(frames.embedding)
).group_by(frames.video_id).collect()

# Identify videos with many scene transitions
dynamic_videos = [v for v in scene_analysis if v['analysis']['scene_changes'] > 5]
print(f"Videos with high scene variability: {len(dynamic_videos)}")
 
```

 
### 3. Custom Quality Metrics for Model Outputs

 
Build domain-specific quality aggregations for [production RAG systems](/blog/production-rag-data-centric):

 
```python

@pxt.uda
class ResponseQualityAggregator(Aggregator):
 """Aggregate quality metrics for LLM responses"""
 
 def __init__(self):
 self.response_lengths = []
 self.confidence_scores = []
 self.error_count = 0
 self.total_count = 0
 
 def update(self, response: str, confidence: float, has_error: bool):
 """Track response quality metrics"""
 self.total_count += 1
 
 if has_error:
 self.error_count += 1
 return
 
 if response:
 self.response_lengths.append(len(response.split()))
 
 if confidence is not None:
 self.confidence_scores.append(confidence)
 
 def value(self) -> dict:
 """Return aggregated quality metrics"""
 return {
 'total_responses': self.total_count,
 'error_rate': self.error_count / self.total_count if self.total_count > 0 else 0,
 'avg_response_length': float(np.mean(self.response_lengths)) if self.response_lengths else 0,
 'avg_confidence': float(np.mean(self.confidence_scores)) if self.confidence_scores else 0,
 'quality_score': (
 (1.0 - (self.error_count / self.total_count)) * 0.5 +
 (np.mean(self.confidence_scores) if self.confidence_scores else 0) * 0.5
 ) if self.total_count > 0 else 0
 }

# Apply to AI responses
from pixeltable.functions import openai

queries = pxt.create_table('customer_queries', {
 'query': pxt.String,
 'customer_id': pxt.String,
 'timestamp': pxt.Timestamp
})

queries.add_computed_column(
 response=openai.chat_completions(
 model='gpt-4o-mini',
 messages=[{'role': 'user', 'content': queries.query}]
 ).choices[0].message.content
)

queries.add_computed_column(
 confidence=0.9 # Placeholder - would come from actual confidence scoring
)

# Aggregate quality by customer
customer_quality = queries.select(
 queries.customer_id,
 quality_metrics=ResponseQualityAggregator(
 queries.response,
 queries.confidence,
 queries.response.errortype.is_not_null()
 )
).group_by(queries.customer_id).collect()

# Identify customers with quality issues
for customer in customer_quality:
 if customer['quality_metrics']['error_rate'] > 0.1:
 print(f"Customer {customer['customer_id']}: {customer['quality_metrics']['error_rate']:.1%} error rate")
 
```

 
## Advanced UDA Patterns for AI Workflows

 
### Rolling Window Aggregations

 
Build time-series aggregations for monitoring AI system performance:

 
```python

@pxt.uda(allows_window=True)
class RollingEmbeddingDrift(Aggregator):
 """Detect embedding drift over time using rolling window"""
 
 def __init__(self, window_size: int = 100):
 self.window_size = window_size
 self.embeddings = []
 self.baseline_embedding = None
 
 def update(self, embedding: np.ndarray, timestamp: float):
 """Track embeddings with timestamps"""
 if embedding is None:
 return
 
 self.embeddings.append((timestamp, embedding))
 
 # Keep only recent window
 if len(self.embeddings) > self.window_size:
 self.embeddings.pop(0)
 
 # Set baseline from first embedding
 if self.baseline_embedding is None:
 self.baseline_embedding = embedding
 
 def value(self) -> dict:
 """Calculate drift from baseline"""
 if not self.embeddings or self.baseline_embedding is None:
 return {'drift_score': 0.0, 'window_size': 0}
 
 # Calculate drift scores
 drift_scores = []
 for timestamp, emb in self.embeddings:
 drift = np.linalg.norm(emb - self.baseline_embedding)
 drift_scores.append(drift)
 
 return {
 'drift_score': float(np.mean(drift_scores)),
 'max_drift': float(np.max(drift_scores)),
 'window_size': len(self.embeddings),
 'is_drifting': float(np.mean(drift_scores)) > 0.5
 }

# Monitor embedding drift over time for model degradation
drift_analysis = frames.select(
 drift_metrics=RollingEmbeddingDrift(
 frames.embedding,
 frames.timestamp
 )
).collect()
 
```

 
### Multimodal Content Aggregations

 
Aggregate across different modalities for comprehensive analytics:

 
```python

@pxt.uda
class MultimodalContentAggregator(Aggregator):
 """Aggregate statistics across multiple modalities"""
 
 def __init__(self):
 self.video_count = 0
 self.image_count = 0
 self.audio_duration = 0.0
 self.total_frames = 0
 self.modalities_present = set()
 
 def update(
 self,
 has_video: bool,
 has_image: bool,
 has_audio: bool,
 audio_duration: float,
 frame_count: int
 ):
 """Track multimodal content statistics"""
 if has_video:
 self.video_count += 1
 self.modalities_present.add('video')
 self.total_frames += frame_count or 0
 
 if has_image:
 self.image_count += 1
 self.modalities_present.add('image')
 
 if has_audio and audio_duration:
 self.audio_duration += audio_duration
 self.modalities_present.add('audio')
 
 def value(self) -> dict:
 """Return comprehensive multimodal statistics"""
 return {
 'video_count': self.video_count,
 'image_count': self.image_count,
 'total_audio_hours': self.audio_duration / 3600,
 'total_frames': self.total_frames,
 'modality_count': len(self.modalities_present),
 'modalities': list(self.modalities_present),
 'is_multimodal': len(self.modalities_present) > 1
 }

# Analyze content library composition
content = pxt.create_table('media_library', {
 'content_id': pxt.String,
 'video': pxt.Video,
 'image': pxt.Image,
 'audio': pxt.Audio,
 'category': pxt.String
})

# Get comprehensive statistics by category
content_stats = content.select(
 content.category,
 stats=MultimodalContentAggregator(
 content.video.is_not_null(),
 content.image.is_not_null(),
 content.audio.is_not_null(),
 content.audio.duration,
 content.frame_count
 )
).group_by(content.category).collect()
 
```

 
## Production Use Cases for Custom Aggregations

 
### Medical Imaging: Diagnostic Confidence Aggregation

 
Build specialized metrics for healthcare AI applications:

 
```python

@pxt.uda
class DiagnosticConfidenceAggregator(Aggregator):
 """Aggregate diagnostic predictions for medical imaging"""
 
 def __init__(self):
 self.predictions = []
 self.slice_count = 0
 self.abnormality_count = 0
 
 def update(self, prediction: dict, slice_location: float):
 """Track predictions across medical imaging slices"""
 if not prediction:
 return
 
 self.slice_count += 1
 
 # Track abnormalities
 if prediction.get('abnormality_detected', False):
 self.abnormality_count += 1
 self.predictions.append({
 'confidence': prediction.get('confidence', 0),
 'location': slice_location,
 'finding': prediction.get('finding', 'unknown')
 })
 
 def value(self) -> dict:
 """Return diagnostic summary"""
 if not self.predictions:
 return {
 'has_findings': False,
 'abnormality_rate': 0.0,
 'avg_confidence': 0.0,
 'critical_slices': []
 }
 
 avg_confidence = sum(p['confidence'] for p in self.predictions) / len(self.predictions)
 
 # Identify critical slices (high confidence abnormalities)
 critical = [
 p for p in self.predictions 
 if p['confidence'] > 0.8
 ]
 
 return {
 'has_findings': True,
 'abnormality_rate': self.abnormality_count / self.slice_count,
 'avg_confidence': avg_confidence,
 'finding_count': len(self.predictions),
 'critical_slices': len(critical),
 'requires_review': avg_confidence > 0.7
 }

# Aggregate medical imaging results by patient
medical_scans = pxt.create_table('radiology.scans', {
 'patient_id': pxt.String,
 'scan_image': pxt.Image,
 'slice_number': pxt.Int,
 'slice_location_mm': pxt.Float
})

# AI predictions on each slice
medical_scans.add_computed_column(
 prediction=analyze_medical_image(medical_scans.scan_image) # Custom AI model
)

# Aggregate by patient for diagnostic summary
patient_diagnostics = medical_scans.select(
 medical_scans.patient_id,
 diagnostic_summary=DiagnosticConfidenceAggregator(
 medical_scans.prediction,
 medical_scans.slice_location_mm
 )
).group_by(medical_scans.patient_id).collect()

# Flag patients needing urgent review
urgent_cases = [
 p for p in patient_diagnostics 
 if p['diagnostic_summary']['requires_review']
]
 
```

 
### Video Processing: Frame Quality Aggregation

 
Monitor video processing quality across [object detection workflows](/blog/object-detection-videos-yolox):

 
```python

@pxt.uda
class VideoQualityAggregator(Aggregator):
 """Aggregate frame quality metrics for video analysis"""
 
 def __init__(self):
 self.blur_scores = []
 self.brightness_scores = []
 self.object_detection_confidences = []
 self.low_quality_frames = []
 
 def update(
 self,
 frame_idx: int,
 blur_score: float,
 brightness: float,
 detections: dict
 ):
 """Track quality metrics per frame"""
 self.blur_scores.append(blur_score)
 self.brightness_scores.append(brightness)
 
 # Track object detection confidence
 if detections and detections.get('boxes'):
 avg_conf = sum(box['confidence'] for box in detections['boxes']) / len(detections['boxes'])
 self.object_detection_confidences.append(avg_conf)
 
 # Flag low quality frames
 if blur_score 215:
 self.low_quality_frames.append(frame_idx)
 
 def value(self) -> dict:
 """Return video quality summary"""
 return {
 'avg_blur': float(np.mean(self.blur_scores)) if self.blur_scores else 0,
 'avg_brightness': float(np.mean(self.brightness_scores)) if self.brightness_scores else 0,
 'avg_detection_confidence': float(np.mean(self.object_detection_confidences)) if self.object_detection_confidences else 0,
 'low_quality_frame_count': len(self.low_quality_frames),
 'low_quality_frame_indices': self.low_quality_frames[:10], # First 10
 'overall_quality_score': self._calculate_overall_quality(),
 'needs_reprocessing': self._calculate_overall_quality() float:
 """Calculate composite quality score"""
 if not self.blur_scores:
 return 0.0
 
 # Normalize metrics
 blur_norm = min(np.mean(self.blur_scores) / 100, 1.0)
 brightness_norm = 1.0 - abs(127.5 - np.mean(self.brightness_scores)) / 127.5
 detection_norm = np.mean(self.object_detection_confidences) if self.object_detection_confidences else 0
 
 return float((blur_norm + brightness_norm + detection_norm) / 3)

# Assess video processing quality
from pixeltable.functions import yolox

frames.add_computed_column(
 detections=yolox(frames.frame, model_id='yolox_s')
)

frames.add_computed_column(
 blur_score=calculate_blur(frames.frame) # Custom blur detection UDF
)

frames.add_computed_column(
 brightness=calculate_brightness(frames.frame) # Custom brightness UDF
)

# Aggregate video quality
video_quality = frames.select(
 frames.video_id,
 quality=VideoQualityAggregator(
 frames.frame_idx,
 frames.blur_score,
 frames.brightness,
 frames.detections
 )
).group_by(frames.video_id).collect()

# Filter videos needing reprocessing
poor_quality_videos = [
 v for v in video_quality 
 if v['quality']['needs_reprocessing']
]
print(f"Videos requiring reprocessing: {len(poor_quality_videos)}")
 
```

 
## Financial Analytics: Custom Portfolio Metrics

 
Build specialized financial aggregations for trading systems:

 
```python

@pxt.uda
class PortfolioRiskAggregator(Aggregator):
 """Calculate portfolio risk metrics"""
 
 def __init__(self):
 self.returns = []
 self.positions = []
 self.volatilities = []
 
 def update(
 self,
 symbol: str,
 position_size: float,
 daily_return: float,
 volatility: float
 ):
 """Accumulate portfolio positions and metrics"""
 self.positions.append({
 'symbol': symbol,
 'size': position_size,
 'return': daily_return,
 'volatility': volatility
 })
 
 self.returns.append(daily_return * position_size)
 self.volatilities.append(volatility)
 
 def value(self) -> dict:
 """Calculate portfolio-level risk metrics"""
 if not self.positions:
 return {'portfolio_return': 0, 'portfolio_risk': 0, 'sharpe_ratio': 0}
 
 # Weighted portfolio return
 portfolio_return = sum(self.returns)
 
 # Portfolio volatility (simplified)
 portfolio_risk = float(np.sqrt(np.mean(np.square(self.volatilities))))
 
 # Sharpe ratio (assuming 3% risk-free rate)
 risk_free_rate = 0.03 / 252 # Daily risk-free rate
 sharpe_ratio = (portfolio_return - risk_free_rate) / portfolio_risk if portfolio_risk > 0 else 0
 
 # Concentration risk
 total_value = sum(abs(p['size']) for p in self.positions)
 max_position = max(abs(p['size']) for p in self.positions) if self.positions else 0
 concentration = max_position / total_value if total_value > 0 else 0
 
 return {
 'portfolio_return': float(portfolio_return),
 'portfolio_risk': portfolio_risk,
 'sharpe_ratio': float(sharpe_ratio),
 'concentration_risk': float(concentration),
 'num_positions': len(self.positions),
 'needs_rebalancing': concentration > 0.25
 }

# Apply to trading data
positions = pxt.create_table('portfolio.positions', {
 'symbol': pxt.String,
 'position_size': pxt.Float,
 'daily_return': pxt.Float,
 'volatility': pxt.Float,
 'date': pxt.Timestamp
})

# Calculate daily portfolio metrics
daily_metrics = positions.select(
 positions.date,
 metrics=PortfolioRiskAggregator(
 positions.symbol,
 positions.position_size,
 positions.daily_return,
 positions.volatility
 )
).group_by(positions.date).collect()
 
```

 
## UDA vs. Traditional Approaches

 
| Approach | Traditional SQL Aggregates | Pixeltable UDAs |
| --- | --- | --- |
| Data Types Supported | Numbers, strings, timestamps only | Embeddings, images, videos, any Python type |
| Custom Logic | Limited to SQL expressions | Full Python with NumPy, scikit-learn, etc. |
| Multimodal Analytics | Requires manual post-processing | Native support for AI data types |
| Integration | Separate tools for complex metrics | Seamless with Pixeltable workflows |
| Performance | Optimized for simple operations | Efficient for complex AI computations |

 
## Best Practices for Building UDAs

 
### Efficient State Management

 

 - **Minimize memory:** Don't store unnecessary data in aggregator state

 - **Use NumPy:** Leverage vectorized operations for performance

 - **Handle nulls:** Always check for None values in update()

 - **Avoid side effects:** UDAs should be pure functions of their input

 

 
### Performance Optimization

 
```python

@pxt.uda
class EfficientAggregator(Aggregator):
 """Example of efficient aggregator design"""
 
 def __init__(self):
 # ✅ Use compact data structures
 self.values = np.array([], dtype=np.float32) # Not list
 self.count = 0
 
 def update(self, value: float):
 # ✅ Handle None early
 if value is None:
 return
 
 # ✅ Efficient append
 self.values = np.append(self.values, value)
 self.count += 1
 
 def value(self) -> float:
 # ✅ Use NumPy for fast computation
 return float(np.mean(self.values)) if self.count > 0 else 0.0
 
```

 
### Debugging and Testing

 
```python

# Test your UDA on small datasets first
test_data = pxt.create_table('test_uda', {
 'value': pxt.Float,
 'category': pxt.String
})

test_data.insert([
 {'value': 1.0, 'category': 'A'},
 {'value': 2.0, 'category': 'A'},
 {'value': 3.0, 'category': 'B'},
])

# Test aggregation
result = test_data.select(
 test_data.category,
 custom_metric=YourCustomUDA(test_data.value)
).group_by(test_data.category).collect()

print("Test results:", result)

# Verify against expected values
assert result[0]['custom_metric'] == expected_value, "UDA logic error"
 
```

 
## Advanced UDA Features

 
### Window Functions with UDAs

 
Enable windowing operations for time-series AI data:

 
```python

@pxt.uda(allows_window=True, requires_order_by=True)
class MovingAverageEmbeddingSimilarity(Aggregator):
 """Calculate moving average of embedding similarities"""
 
 def __init__(self, window_size: int = 10):
 self.window_size = window_size
 self.embedding_window = []
 self.similarity_window = []
 
 def update(self, embedding: np.ndarray):
 """Update with ordered embeddings"""
 if embedding is None:
 return
 
 # Add to window
 self.embedding_window.append(embedding)
 
 # Maintain window size
 if len(self.embedding_window) > self.window_size:
 self.embedding_window.pop(0)
 
 # Calculate similarity with window centroid
 if len(self.embedding_window) > 1:
 centroid = np.mean(self.embedding_window, axis=0)
 similarity = np.dot(embedding, centroid) / (
 np.linalg.norm(embedding) * np.linalg.norm(centroid)
 )
 self.similarity_window.append(similarity)
 
 def value(self) -> float:
 """Return moving average similarity"""
 return float(np.mean(self.similarity_window)) if self.similarity_window else 1.0

# Apply with window clause
windowed_similarity = frames.select(
 frames.frame_idx,
 moving_avg=MovingAverageEmbeddingSimilarity(frames.embedding)
).window(
 order_by=frames.frame_idx,
 rows_between=(10, 0) # 10 frames before to current
).collect()
 
```

 
## Getting Started with Custom Aggregations

 
Ready to build your own AI-specific aggregation functions? Here's how to get started:

 
### Prerequisites

 
```bash

# Install Pixeltable with AI capabilities
pip install pixeltable numpy

# For specific AI providers
pip install pixeltable[openai,huggingface]
 
```

 
### Your First UDA

 
Start with a simple aggregation, then expand to more complex use cases:

 
```python

import pixeltable as pxt
from pixeltable import Aggregator

@pxt.uda
class SimpleCounter(Aggregator):
 """Count non-null values (basic example)"""
 
 def __init__(self):
 self.count = 0
 
 def update(self, value):
 if value is not None:
 self.count += 1
 
 def value(self) -> int:
 return self.count

# Test it
result = table.select(
 total=SimpleCounter(table.column)
).collect()

print(f"Total non-null values: {result[0]['total']}")
 
```

 
## Conclusion: Custom Aggregations for AI-Native Analytics

 
User-Defined Aggregates unlock a new dimension of analytics for AI workflows. Instead of forcing AI data through SQL's limited aggregation functions or resorting to manual post-processing, you can build custom aggregations that understand embeddings, video frames, multimodal content, and domain-specific metrics.

 
 
Combined with Pixeltable's [declarative infrastructure](/blog/declarative-multimodal-incremental), [incremental computation](/blog/incremental-embedding-indexes), and [automatic versioning](/blog/pixeltable-versioning-time-travel), UDAs enable sophisticated analytics that would be impractical with traditional data infrastructure.

 
 
Whether you're building medical imaging analytics, video quality monitoring, financial portfolio analysis, or custom AI metrics, UDAs give you the power to aggregate data exactly how your domain requires, without sacrificing the benefits of declarative data management.

 
## Master Custom Aggregations

 

 - **[UDA SDK Reference](https://docs.pixeltable.com/sdk/latest/pixeltable#uda)**: Official documentation

 - **[Python UDFs Guide](/blog/python-udfs-pixeltable)**: Learn custom functions first

 - **[Video Analytics Tutorial](/blog/object-detection-videos-yolox)**: Apply UDAs to video processing

 - **[Your First Pixeltable Project](/blog/your-first-pixeltable-project)**: Get started with basics

 - **[Pixeltable on GitHub](https://github.com/pixeltable/pixeltable)**: Examples and source code

 - **[Join our Discord](https://discord.gg/QPyqFYx2UN)**: Get help building custom aggregations

 

 
*Stop forcing AI data through SQL's limited aggregations. Build custom metrics that understand your domain.* For a production case study on choosing UDA over UDF when folding video frames into one parent row, see [One Swing, Many Frames: PixelGolf and @pxt.uda](/blog/one-swing-many-frames-pixelgolf-uda-aggregation).