• Home
  • Features
  • Pricing
  • Docs
  • Announcements
  • Sign In

llnl-asr / dftracer-utils / 36502607218

29 Sep 2026 12:20AM UTC coverage: 57.996% (+5.4%) from 52.58%
36502607218

push

github

rayandrew
Merge branch 'feat/view-as-lazyframe' into 'develop'

feat: view as lazyframe

See merge request dftracer/dftracer-utils!25

74077 of 162805 branches covered (45.5%)

Branch coverage included in aggregate %.

3652 of 3978 new or added lines in 49 files covered. (91.8%)

3464 existing lines in 85 files now uncovered.

65902 of 78556 relevant lines covered (83.89%)

179668.29 hits per line

Source File
Press 'n' to go to next uncovered line, 'b' for previous

3.49
/src/dftracer/utils/python/arrow_parallel_reader.cpp
1
#include <dftracer/utils/core/common/config.h>
2
#ifdef DFTRACER_UTILS_ENABLE_ARROW_IPC
3

4
#define PY_SSIZE_T_CLEAN
5
#include <Python.h>
6
#include <dftracer/utils/python/arrow_parallel_reader.h>
7
#include <dftracer/utils/python/py_dict_helpers.h>
8
#include <dftracer/utils/python/py_list_helpers.h>
9
#include <dftracer/utils/python/py_method.h>
10
#include <dftracer/utils/python/py_runtime_mixin.h>
11
#include <dftracer/utils/python/runtime.h>
12
#include <dftracer/utils/python/trace_reader_iterator.h>
13
#include <dftracer/utils/utilities/common/arrow/parallel_reader.h>
14

15
#include <string>
16
#include <vector>
17

18
namespace dftracer::utils::python {
19

20
using utilities::common::arrow::ArrowExportResult;
21
using utilities::common::arrow::read_arrow_files_parallel;
22

UNCOV
23
static PyObject* py_read_arrow_files_parallel(PyObject* /*self*/,
×
24
                                              PyObject* args,
25
                                              PyObject* kwargs) {
26
    static const char* kwlist[] = {"paths", "runtime", nullptr};
27
    PyObject* paths_obj = nullptr;
×
UNCOV
28
    PyObject* runtime_obj = nullptr;
×
29

UNCOV
30
    if (!PyArg_ParseTupleAndKeywords(args, kwargs, "O|O",
×
31
                                     const_cast<char**>(kwlist), &paths_obj,
32
                                     &runtime_obj)) {
UNCOV
33
        return nullptr;
×
34
    }
35

36
    // Convert paths to vector<string>
37
    std::vector<std::string> paths;
×
UNCOV
38
    if (!parse_str_list(paths_obj, "paths", paths)) return nullptr;
×
39

40
    // Get runtime
41
    Runtime* runtime = nullptr;
×
42
    if (runtime_obj && runtime_obj != Py_None) {
×
43
        if (!PyObject_TypeCheck(runtime_obj, &RuntimeType)) {
×
UNCOV
44
            PyErr_SetString(PyExc_TypeError,
×
45
                            "runtime must be a Runtime object");
UNCOV
46
            return nullptr;
×
47
        }
UNCOV
48
        runtime = ((RuntimeObject*)runtime_obj)->runtime.get();
×
49
    } else {
UNCOV
50
        runtime = get_default_runtime();
×
51
    }
52

53
    // Call C++ parallel reader (releases GIL during file I/O)
54
    utilities::common::arrow::ParallelReadResult result;
×
55
    if (!run_blocking_r(
×
56
            [&] {
×
57
                auto task = read_arrow_files_parallel(std::move(paths));
×
58
                return runtime->submit(std::move(task), "read_arrow_files")
×
59
                    .get();
×
UNCOV
60
            },
×
61
            result)) {
UNCOV
62
        return nullptr;
×
63
    }
64

65
    // Build Python result dict
66
    PyObject* file_results_list = PyList_New(result.file_results.size());
×
UNCOV
67
    if (!file_results_list) return nullptr;
×
68

69
    for (std::size_t i = 0; i < result.file_results.size(); ++i) {
×
70
        const auto& fr = result.file_results[i];
×
71
        PyObject* fr_dict = PyDict_New();
×
UNCOV
72
        if (!fr_dict) {
×
73
            Py_DECREF(file_results_list);
×
UNCOV
74
            return nullptr;
×
75
        }
76

77
        dict_set_str(fr_dict, "path", fr.path.c_str());
×
UNCOV
78
        dict_set_bool(fr_dict, "success", fr.success);
×
79

80
        if (!fr.error.empty()) {
×
UNCOV
81
            dict_set_str(fr_dict, "error", fr.error.c_str());
×
82
        } else {
83
            Py_INCREF(Py_None);
×
UNCOV
84
            PyDict_SetItemString(fr_dict, "error", Py_None);
×
85
        }
86

UNCOV
87
        dict_set_i64(fr_dict, "total_rows", fr.total_rows);
×
88

89
        // batches - list of ArrowBatchCapsule objects
90
        PyObject* batches_list = PyList_New(fr.batches->size());
×
UNCOV
91
        if (!batches_list) {
×
92
            Py_DECREF(fr_dict);
×
93
            Py_DECREF(file_results_list);
×
UNCOV
94
            return nullptr;
×
95
        }
96

UNCOV
97
        for (std::size_t j = 0; j < fr.batches->size(); ++j) {
×
98
            ArrowBatchCapsuleObject* capsule =
UNCOV
99
                (ArrowBatchCapsuleObject*)ArrowBatchCapsuleType.tp_alloc(
×
100
                    &ArrowBatchCapsuleType, 0);
UNCOV
101
            if (!capsule) {
×
102
                Py_DECREF(batches_list);
×
103
                Py_DECREF(fr_dict);
×
104
                Py_DECREF(file_results_list);
×
UNCOV
105
                return nullptr;
×
106
            }
107
            // Move the batch into the capsule
108
            capsule->result =
×
109
                new ArrowExportResult(std::move((*fr.batches)[j]));
×
UNCOV
110
            PyList_SetItem(batches_list, j, (PyObject*)capsule);
×
111
        }
112

UNCOV
113
        PyDict_SetItemString(fr_dict, "batches", batches_list);
×
114
        Py_DECREF(batches_list);
×
115

UNCOV
116
        PyList_SetItem(file_results_list, i, fr_dict);
×
117
    }
118

119
    // Build final result dict
120
    PyObject* result_dict = PyDict_New();
×
UNCOV
121
    if (!result_dict) {
×
122
        Py_DECREF(file_results_list);
×
UNCOV
123
        return nullptr;
×
124
    }
125

UNCOV
126
    PyDict_SetItemString(result_dict, "file_results", file_results_list);
×
127
    Py_DECREF(file_results_list);
×
128

129
    dict_set_i64(result_dict, "total_rows", result.total_rows);
×
130
    dict_set_i64(result_dict, "total_batches", result.total_batches);
×
131
    dict_set_size(result_dict, "files_read", result.files_read);
×
UNCOV
132
    dict_set_size(result_dict, "files_failed", result.files_failed);
×
133

134
    return result_dict;
×
UNCOV
135
}
×
136

137
static PyMethodDef arrow_parallel_reader_methods[] = {
138
    {"read_arrow_files_parallel",
139
     DFTU_PYCFUNCTION(py_read_arrow_files_parallel),
140
     METH_VARARGS | METH_KEYWORDS,
141
     "Read multiple Arrow IPC files in parallel using the Runtime.\n\n"
142
     "Args:\n"
143
     "    paths: List of file paths to read.\n"
144
     "    runtime: Optional Runtime object. Uses default if not provided.\n\n"
145
     "Returns:\n"
146
     "    dict with:\n"
147
     "        - file_results: List of per-file results, each with:\n"
148
     "            - path: File path\n"
149
     "            - success: True if read succeeded\n"
150
     "            - error: Error message if failed, else None\n"
151
     "            - total_rows: Number of rows in file\n"
152
     "            - batches: List of ArrowBatch objects\n"
153
     "        - total_rows: Total rows across all files\n"
154
     "        - total_batches: Total batches across all files\n"
155
     "        - files_read: Number of files read successfully\n"
156
     "        - files_failed: Number of files that failed"},
157
    {nullptr, nullptr, 0, nullptr}};
158

159
int init_arrow_parallel_reader(PyObject* m) {
2 ✔
160
    // Add the function to the module
161
    if (PyModule_AddFunctions(m, arrow_parallel_reader_methods) < 0) {
2 ✔
UNCOV
162
        return -1;
×
163
    }
164
    return 0;
2 ✔
165
}
1 ✔
166

167
}  // namespace dftracer::utils::python
168

169
#endif  // DFTRACER_UTILS_ENABLE_ARROW_IPC
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE TRIAL · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc