Implementation:NVIDIA DALI C API V2 Pipeline
| Knowledge Sources | |
|---|---|
| Domains | Data_Pipeline, C_API |
| Last Updated | 2026-02-08 16:00 GMT |
Overview
The C API v2 pipeline implementation provides C-callable functions for the full DALI pipeline lifecycle including creation, deserialization, building, running, input feeding, output retrieval, and checkpointing.
Description
This file implements the C API v2 pipeline management functions that operate through the PipelineWrapper class defined in the corresponding header. It belongs to the dali::c_api namespace within NVIDIA DALI. The implementation provides a modernized replacement for the legacy C API with structured error handling (returning daliResult_t codes), typed parameter structs, and proper separation of pipeline outputs from the pipeline handle.
The file contains three main sections: the PipelineWrapper C++ method implementations, the CheckpointWrapper serialization logic, and the C function entry points. The ToCppParams helper converts the C-style daliPipelineParams_t struct (with per-field presence flags) into the C++ PipelineParams struct. The wrapper methods delegate to the underlying dali::Pipeline while adding input/output descriptor introspection, feed-input validation, and checkpoint management.
Key C functions include: daliPipelineCreate and daliPipelineDeserialize for pipeline construction, daliPipelineBuild/daliPipelineRun/daliPipelinePrefetch for execution, daliPipelineFeedInput for external data ingestion with configurable copy modes, daliPipelinePopOutputs/daliPipelinePopOutputsAsync for output retrieval, and a complete checkpoint API (daliPipelineGetCheckpoint, daliPipelineSerializeCheckpoint, daliPipelineDeserializeCheckpoint, daliPipelineRestoreCheckpoint).
Usage
Use these functions to manage DALI pipelines from C code or language bindings. Create pipelines from parameters or by deserializing protobuf, build and run them, feed external inputs with optional copy control, pop outputs as independent objects, and use checkpointing for fault-tolerant training workflows.
Code Reference
Source Location
- Repository: NVIDIA_DALI
- File: dali/c_api_2/pipeline.cc
- Lines: 1-515
Signature
// Pipeline lifecycle
daliResult_t daliPipelineCreate(daliPipeline_h *out_pipe_handle,
const daliPipelineParams_t *params);
daliResult_t daliPipelineDeserialize(daliPipeline_h *out_pipe_handle,
const void *serialized_pipeline,
size_t serialized_pipeline_size,
const daliPipelineParams_t *param_overrides);
daliResult_t daliPipelineBuild(daliPipeline_h pipeline);
daliResult_t daliPipelineRun(daliPipeline_h pipeline);
daliResult_t daliPipelinePrefetch(daliPipeline_h pipeline);
daliResult_t daliPipelineDestroy(daliPipeline_h pipeline);
// Input/Output
daliResult_t daliPipelineFeedInput(daliPipeline_h pipeline, const char *input_name,
daliTensorList_h input_data, const char *data_id,
daliFeedInputFlags_t options, const cudaStream_t *stream);
daliResult_t daliPipelinePopOutputs(daliPipeline_h pipeline, daliPipelineOutputs_h *out);
daliResult_t daliPipelinePopOutputsAsync(daliPipeline_h pipeline, daliPipelineOutputs_h *out,
cudaStream_t stream);
daliResult_t daliPipelineGetFeedCount(daliPipeline_h pipeline, int *out_feed_count,
const char *input_name);
// Introspection
daliResult_t daliPipelineGetInputCount(daliPipeline_h pipeline, int *out_input_count);
daliResult_t daliPipelineGetInputDescByIdx(daliPipeline_h pipeline,
daliPipelineIODesc_t *out_input_desc, int index);
daliResult_t daliPipelineGetInputDesc(daliPipeline_h pipeline,
daliPipelineIODesc_t *out_input_desc, const char *name);
daliResult_t daliPipelineGetOutputCount(daliPipeline_h pipeline, int *out_count);
daliResult_t daliPipelineGetOutputDesc(daliPipeline_h pipeline,
daliPipelineIODesc_t *out_desc, int index);
// Checkpointing
daliResult_t daliPipelineGetCheckpoint(daliPipeline_h pipeline, daliCheckpoint_h *out_checkpoint,
const daliCheckpointExternalData_t *checkpoint_ext);
daliResult_t daliPipelineSerializeCheckpoint(daliPipeline_h pipeline, daliCheckpoint_h checkpoint,
const char **out_data, size_t *out_size);
daliResult_t daliPipelineDeserializeCheckpoint(daliPipeline_h pipeline,
daliCheckpoint_h *out_checkpoint,
const char *serialized, size_t size);
daliResult_t daliPipelineRestoreCheckpoint(daliPipeline_h pipeline, daliCheckpoint_h checkpoint);
daliResult_t daliCheckpointDestroy(daliCheckpoint_h checkpoint);
Import
#include "dali/c_api_2/pipeline.h"
I/O Contract
Inputs
| Name | Type | Required | Description |
|---|---|---|---|
| params | const daliPipelineParams_t * | Yes | Pipeline parameters with presence flags for optional fields |
| serialized_pipeline | const void * | Yes (deserialize) | Protobuf-serialized pipeline definition |
| input_name | const char * | Yes | Name of the input operator to feed data to |
| input_data | daliTensorList_h | Yes | TensorList handle containing input data |
| options | daliFeedInputFlags_t | No | Flags: DALI_FEED_INPUT_FORCE_COPY, DALI_FEED_INPUT_NO_COPY, DALI_FEED_INPUT_SYNC |
| stream | const cudaStream_t * | No | CUDA stream for async operations (NULL for default) |
Outputs
| Name | Type | Description |
|---|---|---|
| daliResult_t | enum | Operation result code |
| out_pipe_handle | daliPipeline_h * | Created pipeline handle |
| out | daliPipelineOutputs_h * | Pipeline outputs handle (caller must destroy) |
| out_desc | daliPipelineIODesc_t * | I/O descriptor with name, device, dtype, ndim, layout |
| out_checkpoint | daliCheckpoint_h * | Checkpoint handle for serialization/restoration |
Usage Examples
Deserialize, Build, and Run a Pipeline
#include "dali/dali.h"
daliPipelineParams_t params{};
params.max_batch_size_present = true;
params.max_batch_size = 8;
params.prefetch_queue_depths_present = true;
params.prefetch_queue_depths = {3, 3};
daliPipeline_h pipeline = nullptr;
daliResult_t r = daliPipelineDeserialize(&pipeline, proto.data(), proto.size(), ¶ms);
assert(r == DALI_SUCCESS);
daliPipelineBuild(pipeline);
daliPipelinePrefetch(pipeline);
// Get outputs
daliPipelineOutputs_h outputs = nullptr;
daliPipelinePopOutputs(pipeline, &outputs);
// Access output tensor list
daliTensorList_h tl = nullptr;
daliPipelineOutputsGet(outputs, &tl, 0);
// Cleanup
daliTensorListDecRef(tl, nullptr);
daliPipelineOutputsDestroy(outputs);
daliPipelineDestroy(pipeline);
Feed External Input
int feed_count = 0;
daliPipelineGetFeedCount(pipeline, &feed_count, "external_input");
for (int i = 0; i < feed_count; i++) {
cudaStream_t stream = my_stream;
daliPipelineFeedInput(pipeline, "external_input", tensor_list_handle,
"data_id_001", DALI_FEED_INPUT_SYNC, &stream);
}
daliPipelinePrefetch(pipeline);