BaseIngestor
BaseIngestor is the foundation every plugin is built on. The templates you
normally use — GenericZarrIngestor, GenericParquetIngestor,
and DirectZarrIngestor — are all BaseIngestor subclasses. It defines
the engine contract and owns the run() lifecycle; each template fills that
lifecycle in with batteries (file discovery, batching, a write strategy) and
exposes a smaller, friendlier hook in place of the raw ones.
You subclass BaseIngestor directly only when none of the templates fit and
you want to own batch processing yourself. Scaffold one with
firecube plugins create my-plugin --template base (see
Create a Plugin).
How the pieces fit
classDiagram
class BaseIngestor {
<>
+run(ctx)
+_process_batch(batch, ctx)*
+_aggregate_metrics(ctx, state)*
}
class GenericZarrIngestor {
+build_dataset(group, items, ctx)
}
class GenericParquetIngestor {
+build_dataset(group, batch, ctx)
}
class DirectZarrIngestor {
+zarr_schema(ctx)
+build_write_intents(batch, ctx)
}
class YourPlugin
class DuckDbMixin {
<>
}
class PluginContext {
storage
telemetry
options
}
BaseIngestor <|-- GenericZarrIngestor
BaseIngestor <|-- GenericParquetIngestor
BaseIngestor <|-- DirectZarrIngestor
GenericZarrIngestor <|-- YourPlugin
DuckDbMixin <|.. YourPlugin : optional
BaseIngestor ..> PluginContext : passes to hooks
A concrete plugin extends one template (or BaseIngestor directly), optionally
mixes in a helper (see Mixins), and receives a read-only
PluginContext in every hook. The three Generic* templates also share batch-lifecycle hooks
(batch_setup, prepare_batch_data, cleanup_batch_data, batch_teardown);
DirectZarrIngestor dispatches WriteIntents instead.
What the engine runs
When a run starts, BaseIngestor.run() drives the whole pipeline and calls your
code at two points. Around your hooks it:
- validates and merges the configuration tiers,
- discovers and batches the source files,
- opens telemetry and acquires the ChunkManager write claims,
- calls
_process_batchonce per batch — your code, - calls
_aggregate_metricsto roll per-batch metrics into the run result — your code, - records the run's outcome so it is idempotent and resumable.
The templates implement steps 4–5 for you (delegating to build_dataset /
build_write_intents); a direct BaseIngestor subclass implements them itself.
Required hooks
| Hook | Signature | Purpose |
|---|---|---|
_process_batch |
(batch, ctx) -> PipelineResult |
Process one batch and return a typed result. |
_aggregate_metrics |
(ctx, state) -> Mapping[str, Any] |
Roll per-batch metrics into the run result. |
For the standard aggregation behavior, call merge_batch_metrics(ctx, state)
from firecube.ingestor.api.
Context contract
- Plugin-facing hooks receive a read-only
PluginContext. - Runtime finalization hooks receive the engine-owned runtime context.
- Legacy
ctx: IngestContextannotations are rejected at runtime with aConfigurationError.
Minimal example
from collections.abc import Mapping
from typing import Any, ClassVar
from firecube.ingestor.api import (
BaseIngestor,
OutputPaths,
PipelineBatch,
PipelineResult,
PluginContext,
merge_batch_metrics,
register_ingestor,
)
@register_ingestor("my_custom_plugin")
class MyCustomIngestor(BaseIngestor):
PRODUCT_NAME: ClassVar[str] = "my_custom_product"
name = "my_custom_plugin"
def _process_batch(
self,
batch: PipelineBatch,
ctx: PluginContext,
) -> PipelineResult:
...
return PipelineResult(
batch=batch,
success=True,
outputs=OutputPaths(primary=str(ctx.target)),
)
def _aggregate_metrics(self, ctx, state) -> Mapping[str, Any]:
return merge_batch_metrics(ctx, state)
Next Steps
- Plugin Contract — check
PRODUCT_NAME, typed outputs, and metric requirements - GenericZarrIngestor — use the append-Zarr template
- GenericParquetIngestor — use the Parquet template
- DirectZarrIngestor — write fixed Zarr regions