This document defines how node metadata is aggregated across batch executions in micro-batching scenarios. It serves as the authoritative reference for understanding, maintaining, and extending metadata aggregation behavior.
This is the primary maintainer guide for updating aggregation behavior when operators introduce new metadata fields.
Each batch execution produces a NodeMetadataItem with the following structure:
{
"id": "node-uuid", # Node identifier
"operator": "Operator Name", # Operator name
"node_metadata": { # Nested operator-specific data
# Base metadata fields (present in all operators)
"total_docs": 100,
"processed_docs": 100,
"failed_docs_count": 0,
"failed_docs": [],
"skipped_docs_count": 0,
"skipped_docs": [],
"node_status": "COMPLETED",
# Operator-specific fields (varies by operator)
"progress_percentage": "100.00%",
"custom_field": "operator-specific value",
# ... other operator-specific fields
}
}
NodeMetadataItem recordbatch_id and batch_num fields| Strategy | Description | Example |
|---|---|---|
| SUM | Add numeric values | 10 + 20 = 30 |
| UNION | Combine lists, remove duplicates | [a,b] + [b,c] = [a,b,c] |
| CONCAT | Combine lists, keep duplicates | [a,b] + [b,c] = [a,b,b,c] |
| PRIORITY_STATUS | Select most severe status | RUNNING > FAILED > COMPLETED |
| MIN | Select minimum value | min(10, 5, 15) = 5 |
| MAX | Select maximum value | max(10, 5, 15) = 15 |
| WEIGHTED_AVERAGE | Average with batch size consideration | (50*10 + 75*20) / 30 = 66.67 |
| LAST_COMPLETED | Use last non-empty value | Last error message |
| MERGE_DICT | Shallow merge dictionaries | {a:1} + {b:2} = {a:1, b:2} |
| DEEP_MERGE | Deep merge nested dictionaries | Recursive merge |
| FIRST | Use first value | First batch value |
| LAST | Use last value | Last batch value |
| CUSTOM | Custom function | User-defined logic |
Strategies are implemented in:
src/docpipe/core/job_management/application/aggregation/strategies.pyAggregationStrategyMetadataAggregator classThese strategies apply to the nested node_metadata field within NodeStatsDto:
DEFAULT_STRATEGIES = {
# Document counters - sum across batches
"processed_docs": AggregationStrategy.SUM,
"total_docs": AggregationStrategy.SUM,
"failed_docs_count": AggregationStrategy.SUM,
"skipped_docs_count": AggregationStrategy.SUM,
# Document ID lists - union to avoid duplicates
"docs_completed": AggregationStrategy.UNION,
"failed_docs": AggregationStrategy.UNION,
"skipped_docs": AggregationStrategy.UNION,
# Status fields - priority-based (RUNNING > FAILED > COMPLETED)
"status": AggregationStrategy.PRIORITY_STATUS,
"node_status": AggregationStrategy.PRIORITY_STATUS,
# Timing fields
"start_time": AggregationStrategy.MIN, # Earliest start
"end_time": AggregationStrategy.MAX, # Latest end
"time_taken": AggregationStrategy.SUM, # Total duration
# Progress tracking
"progress_percentage": AggregationStrategy.WEIGHTED_AVERAGE,
# Error handling - keep last meaningful error
"error": AggregationStrategy.LAST_COMPLETED,
"error_message": AggregationStrategy.LAST_COMPLETED,
}
These fields are aggregated at the NodeStatsDto level (not in nested metadata):
| Field | Strategy | Rationale |
|---|---|---|
node_status |
PRIORITY_STATUS | Most severe status wins |
start_time |
MIN | Earliest batch start |
end_time |
MAX | Latest batch end |
time_taken |
SUM | Total processing time |
total_docs |
UNION | All unique document IDs |
docs_completed |
UNION | All completed document IDs |
failed_docs |
UNION | All failed document IDs |
skipped_docs |
UNION | All skipped document IDs |
docs_completed_count |
Calculated | len(docs_completed) |
col_names |
FIRST | Schema from first batch |
error |
LAST_COMPLETED | Last meaningful error |
node_metadata |
MetadataAggregator | Nested aggregation |
All operators MUST emit these base metadata fields:
from common.constants.constants import Metrics, ExecutionStatus
metadata = {
Metrics.External.TOTAL_DOCS: 100, # Total documents to process
Metrics.External.PROCESSED_DOCS: 95, # Documents processed so far
Metrics.External.FAILED_DOCS_COUNT: 3, # Count of failed documents
Metrics.External.FAILED_DOCS: [ # List of failed documents
{"id": "doc1", "name": "file1.pdf", "reason": "Parse error", "document_url": ""}
],
Metrics.External.SKIPPED_DOCS_COUNT: 2, # Count of skipped documents
Metrics.External.SKIPPED_DOCS: [ # List of skipped documents
{"id": "doc2", "name": "file2.pdf", "reason": "Empty file", "document_url": ""}
],
Metrics.External.NODE_STATUS: ExecutionStatus.COMPLETED.value # Final status
}
Use AbstractOperator helper methods for consistency:
from core.operators.abstract_operator import AbstractOperator
# Create base metadata
metadata = AbstractOperator.create_base_metadata(
total_docs_count=table.num_rows,
node_status=ExecutionStatus.RUNNING.value
)
# Record failed document
AbstractOperator.record_failed_document(
metadata=metadata,
doc_id="doc123",
doc_name="file.pdf",
reason="Processing error"
)
# Record skipped document
AbstractOperator.record_skipped_document(
metadata=metadata,
doc_id="doc456",
doc_name="empty.pdf",
reason="Empty file"
)
Operators MAY add custom fields on top of base metadata:
# Example: Embedding operator
metadata["embedding_model"] = "nomic-embed-text"
metadata["embedding_dimension"] = 768
metadata["batch_size"] = 32
# Example: Extraction operator
metadata["extraction_method"] = "docling"
metadata["tables_extracted"] = 5
metadata["images_extracted"] = 12
When an operator adds, removes, renames, or changes the meaning of any emitted metadata field, maintainers must review aggregation behavior in strategies.py.
LAST behavior is correct for each field.DEFAULT_STRATEGIES.Micro-batch execution stores raw node stats per batch and combines them later during read-side aggregation. If a new metadata field is introduced without an explicit aggregation review, the aggregated API view may silently produce incorrect values.
Common examples:
SUMUNIONWEIGHTED_AVERAGEPRIORITY_STATUSLASTIf a field is not listed in DEFAULT_STRATEGIES, the aggregator falls back to the default LAST strategy. That fallback is intentional, but it is not correct for many counters, lists, or status fields.
When creating a new operator that emits metadata:
class MyNewOperator(AbstractOperator):
def transform(self, table: pa.Table, file_name: str = None) -> tuple[list[pa.Table], dict[str, Any]]:
# 1. Create base metadata
metadata = self.create_base_metadata(
total_docs_count=table.num_rows,
node_status=ExecutionStatus.RUNNING.value
)
# 2. Add operator-specific fields
metadata["my_custom_metric"] = 0
metadata["my_custom_list"] = []
# 3. Process and update metadata
# ... processing logic ...
return [output_table], metadata
This step is mandatory whenever the operator emits new metadata.
strategies.py.LAST behavior.DEFAULT_STRATEGIES for every field that requires explicit aggregation.DEFAULT_STRATEGIES = {
# ... existing strategies ...
"my_custom_metric": AggregationStrategy.SUM,
"my_custom_list": AggregationStrategy.UNION,
}
If using default strategy (LAST):
DEFAULT_STRATEGIESCreate tests for metadata aggregation:
def test_my_operator_metadata_aggregation():
"""Test metadata aggregation across batches."""
# Create batch records
batch1_metadata = {"my_custom_metric": 10, "my_custom_list": ["a", "b"]}
batch2_metadata = {"my_custom_metric": 20, "my_custom_list": ["b", "c"]}
# Aggregate
aggregator = MetadataAggregator()
result = aggregator.aggregate_metadata(
metadata_list=[batch1_metadata, batch2_metadata]
)
# Verify
assert result["my_custom_metric"] == 30 # SUM
assert set(result["my_custom_list"]) == {"a", "b", "c"} # UNION
When modifying operator metadata:
Questions to answer:
If adding new fields:
LAST behavior without reviewing the field explicitlyIf modifying existing fields:
DEFAULT_STRATEGIESRequired test updates:
Required documentation updates:
Location: tests/unit/core/job_management/application/aggregation/
Required tests:
Example:
def test_sum_strategy():
aggregator = MetadataAggregator()
result = aggregator.aggregate_metadata(
metadata_list=[
{"count": 10},
{"count": 20},
{"count": 30}
]
)
assert result["count"] == 60
Location: tests/integration/core/job_management/
Required tests:
Example:
def test_batch_aggregation_with_real_operator():
# Execute operator across 3 batches
batch_records = execute_operator_in_batches(
operator=ExtractDocling(),
batches=3
)
# Aggregate
aggregated = aggregate_batch_node_stats(
node_id="extract_1",
batch_records=batch_records,
aggregator=MetadataAggregator()
)
# Verify aggregation
assert aggregated["node_status"] == "COMPLETED"
assert len(aggregated["total_docs"]) == total_expected_docs
| File | Purpose |
|---|---|
strategies.py |
Strategy enum and default mappings |
aggregator.py |
MetadataAggregator implementation |
batch_aggregator.py |
Batch-level aggregation logic |
node_stats_aggregator.py |
Service layer aggregation |
node_stats_dto.py |
NodeStatsDto and NodeMetadataItem models |
Operator Execution (Batch 1)
↓
NodeMetadataItem (Batch 1) → Store
↓
Operator Execution (Batch 2)
↓
NodeMetadataItem (Batch 2) → Store
↓
Operator Execution (Batch 3)
↓
NodeMetadataItem (Batch 3) → Store
↓
API Request (Get Status)
↓
NodeStatsAggregator.get_aggregated_node_stats()
↓
Fetch all batch records from store
↓
Group by node_id
↓
aggregate_batch_node_stats() for each node
↓
MetadataAggregator.aggregate_metadata()
↓
Apply strategy for each field
↓
Return aggregated NodeStatsDto
This implementation aligns with enterprise behavior documented in:
JOBSTATS_MANAGEMENT_KNOWLEDGE.md Section “Aggregation Strategy”JOBSTATS_MANAGEMENT_KNOWLEDGE.md Section “Node Metadata vs Node Stats”JOBSTATS_MANAGEMENT_KNOWLEDGE.md Section “Micro-Batching Support”Key alignment points:
When working with node metadata:
DEFAULT_STRATEGIESFor questions about metadata aggregation:
JOBSTATS_MANAGEMENT_KNOWLEDGE.mdtests/unit/core/job_management/Document Version: 1.0
Last Updated: 2026-04-21
Maintained By: Docling Pipelines Core Team