• 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

53.15
/src/dftracer/utils/python/lazyframe.cpp
1
// A native LazyFrame exposed to Python: a deferred query over an in-memory
2
// dataframe::DataFrame. Builder methods record ops and return a new LazyFrame;
3
// nothing runs until collect(), which materializes a native _DataFrame. A
4
// filter/with_column expr's col(i) refers to the i-th column of the frame at
5
// that point, so the Python wrapper resolves names to indices against schema().
6

7
#include <dftracer/utils/core/common/config.h>  // DFTRACER_UTILS_ENABLE_ARROW
8
#include <dftracer/utils/python/lazyframe.h>
9

10
#ifdef DFTRACER_UTILS_ENABLE_ARROW
11

12
#include <dftracer/utils/core/common/memory_budget.h>  // NO_SPILL_BUDGET
13
#include <dftracer/utils/core/runtime.h>
14
#include <dftracer/utils/core/tasks/coro_scope.h>
15
#include <dftracer/utils/dataframe/dataframe.h>
16
#include <dftracer/utils/dataframe/expr.h>
17
#include <dftracer/utils/dataframe/lazyframe.h>
18
#include <dftracer/utils/python/columnar_eval.h>
19
#include <dftracer/utils/python/dataframe.h>
20
#include <dftracer/utils/python/py_agg_helpers.h>
21
#include <dftracer/utils/python/py_frame_op_helpers.h>
22
#include <dftracer/utils/python/py_join_helpers.h>
23
#include <dftracer/utils/python/py_list_helpers.h>
24
#include <dftracer/utils/python/py_method.h>
25
#include <dftracer/utils/python/py_runtime_mixin.h>  // run_blocking
26
#include <dftracer/utils/python/py_scalar_helpers.h>
27
#include <dftracer/utils/python/py_seq_helpers.h>
28
#include <dftracer/utils/python/py_type_helpers.h>
29
#include <dftracer/utils/python/series.h>
30
#include <dftracer/utils/python/streaming_iterator.h>
31

32
#include <cstdint>
33
#include <new>
34
#include <stdexcept>
35
#include <string>
36
#include <utility>
37
#include <vector>
38

39
namespace dataframe = dftracer::utils::dataframe;
40

41
namespace {
42

43
using dataframe::GroupAgg;
44
using dataframe::LazyFrame;
45
using dftracer::utils::python::aggs_from_seq;
46
using dftracer::utils::python::group_agg_from_spec;
47
using dftracer::utils::python::parse_string_seq;
48
using dftracer::utils::python::strings_from_str_or_seq;
49

50
struct LazyFrameObject {
51
    PyObject_HEAD LazyFrame lf;
52
};
53

54
PyTypeObject LazyFrameType;
55

56
LazyFrameObject* as_lazyframe(PyObject* o) {
5,012 ✔
57
    if (!PyObject_TypeCheck(o, &LazyFrameType)) {
5,012 ✔
58
        PyErr_SetString(PyExc_TypeError, "expected a LazyFrame");
2 ✔
59
        return nullptr;
2 ✔
60
    }
61
    return reinterpret_cast<LazyFrameObject*>(o);
5,010 ✔
62
}
2,506 ✔
63

64
PyObject* make_lazyframe(LazyFrame&& lf) {
4,466 ✔
65
    auto* self = reinterpret_cast<LazyFrameObject*>(
2,233 ✔
66
        LazyFrameType.tp_alloc(&LazyFrameType, 0));
4,466 ✔
67
    if (!self) return nullptr;
4,466 ✔
68
    new (&self->lf) LazyFrame(std::move(lf));
4,466 ✔
69
    return reinterpret_cast<PyObject*>(self);
4,466 ✔
70
}
2,233 ✔
71

72
void LazyFrame_dealloc(LazyFrameObject* self) {
4,466 ✔
73
    self->lf.~LazyFrame();
4,466 ✔
74
    Py_TYPE(self)->tp_free(reinterpret_cast<PyObject*>(self));
4,466 ✔
75
}
4,466 ✔
76

77
// Wrap a builder call that returns a LazyFrame, turning C++ exceptions into
78
// Python errors.
79
template <class Fn>
80
PyObject* run_lazy_op(Fn&& fn) {
2,374 ✔
81
    try {
82
        return make_lazyframe(fn());
2,374 !
83
    } catch (const std::exception& e) {
6 !
84
        PyErr_SetString(PyExc_ValueError, e.what());
6 !
85
        return nullptr;
6 ✔
86
    }
3 !
87
}
1,190 ✔
88

89
// filter(ast) -> LazyFrame. `ast` is the columnar.py post-order Expr AST; the
90
// engine rebuilds a dataframe::Expr referencing columns by position.
91
PyObject* LazyFrame_filter(PyObject* self, PyObject* ast) {
620 ✔
92
    LazyFrameObject* b = as_lazyframe(self);
620 !
93
    if (!b) return nullptr;
620 ✔
94
    dataframe::Expr pred;
620 ✔
95
    if (!dftracer::utils::python::build_expr_from_ast(ast, &pred))
620 !
96
        return nullptr;
×
97
    return run_lazy_op([&] { return b->lf.filter(pred); });
1,240 !
98
}
620 ✔
99

100
PyObject* LazyFrame_with_column(PyObject* self, PyObject* args) {
532 ✔
101
    LazyFrameObject* b = as_lazyframe(self);
532 !
102
    if (!b) return nullptr;
532 ✔
103
    const char* name = nullptr;
532 ✔
104
    PyObject* ast = nullptr;
532 ✔
105
    if (!PyArg_ParseTuple(args, "sO", &name, &ast)) return nullptr;
532 !
106
    dataframe::Expr e;
532 ✔
107
    if (!dftracer::utils::python::build_expr_from_ast(ast, &e)) return nullptr;
532 !
108
    return run_lazy_op([&] { return b->lf.with_column(name, e); });
1,064 !
109
}
532 ✔
110

111
PyObject* LazyFrame_select(PyObject* self, PyObject* names) {
396 ✔
112
    LazyFrameObject* b = as_lazyframe(self);
396 !
113
    if (!b) return nullptr;
396 ✔
114
    std::vector<std::string> ns;
396 ✔
115
    if (!parse_str_list(names, "names", ns)) return nullptr;
396 !
116
    return run_lazy_op([&] { return b->lf.select(std::move(ns)); });
792 !
117
}
396 ✔
118

119
PyObject* LazyFrame_rename(PyObject* self, PyObject* names) {
12 ✔
120
    LazyFrameObject* b = as_lazyframe(self);
12 !
121
    if (!b) return nullptr;
12 ✔
122
    std::vector<std::string> ns;
12 ✔
123
    if (!parse_str_list(names, "names", ns)) return nullptr;
12 !
124
    return run_lazy_op([&] { return b->lf.rename(std::move(ns)); });
24 !
125
}
12 ✔
126

127
PyObject* LazyFrame_slice(PyObject* self, PyObject* args) {
14 ✔
128
    LazyFrameObject* b = as_lazyframe(self);
14 !
129
    if (!b) return nullptr;
14 ✔
130
    long long offset = 0, len = 0;
14 ✔
131
    if (!PyArg_ParseTuple(args, "LL", &offset, &len)) return nullptr;
14 !
132
    return run_lazy_op([&] { return b->lf.slice(offset, len); });
28 !
133
}
7 ✔
134

135
PyObject* LazyFrame_head(PyObject* self, PyObject* n) {
10 ✔
136
    LazyFrameObject* b = as_lazyframe(self);
10 !
137
    if (!b) return nullptr;
10 ✔
138
    long long v = PyLong_AsLongLong(n);
10 !
139
    if (v == -1 && PyErr_Occurred()) return nullptr;
10 !
140
    return run_lazy_op([&] { return b->lf.head(v); });
20 !
141
}
5 ✔
142

143
PyObject* LazyFrame_take(PyObject* self, PyObject* seq) {
8 ✔
144
    LazyFrameObject* b = as_lazyframe(self);
8 !
145
    if (!b) return nullptr;
8 ✔
146
    std::vector<std::int64_t> idx;
8 ✔
147
    if (!dftracer::utils::python::parse_int_seq(
8 !
148
            seq, "take() expects a sequence of ints", idx))
4 ✔
149
        return nullptr;
×
150
    return run_lazy_op([&] { return b->lf.take(std::move(idx)); });
16 !
151
}
8 ✔
152

153
PyObject* LazyFrame_filter_mask(PyObject* self, PyObject* mask) {
8 ✔
154
    LazyFrameObject* b = as_lazyframe(self);
8 !
155
    if (!b) return nullptr;
8 ✔
156
    const dataframe::Series* m =
4 ✔
157
        dftracer::utils::python::unwrap_vec_column(mask);
8 !
158
    if (!m) {
8 ✔
159
        PyErr_SetString(PyExc_TypeError,
×
160
                        "filter_mask() expects a Series boolean mask");
161
        return nullptr;
×
162
    }
163
    return run_lazy_op([&] { return b->lf.filter_mask(m->share()); });
20 !
164
}
4 ✔
165

166
PyObject* LazyFrame_reverse(PyObject* self, PyObject*) {
4 ✔
167
    LazyFrameObject* b = as_lazyframe(self);
4 !
168
    if (!b) return nullptr;
4 ✔
169
    return run_lazy_op([&] { return b->lf.reverse(); });
8 !
170
}
2 ✔
171

172
PyObject* LazyFrame_sort_by_multi(PyObject* self, PyObject* args,
80 ✔
173
                                  PyObject* kwds) {
174
    LazyFrameObject* b = as_lazyframe(self);
80 !
175
    if (!b) return nullptr;
80 ✔
176
    PyObject* names_obj = nullptr;
80 ✔
177
    PyObject* descending_obj = nullptr;
80 ✔
178
    static const char* kw[] = {"names", "descending", nullptr};
179
    if (!PyArg_ParseTupleAndKeywords(args, kwds, "O|O", const_cast<char**>(kw),
80 !
180
                                     &names_obj, &descending_obj))
181
        return nullptr;
×
182
    std::vector<std::string> names;
80 ✔
183
    if (!strings_from_str_or_seq(names_obj, names)) return nullptr;
80 !
184
    std::vector<bool> descending;
80 ✔
185
    if (!descending_obj) {
80 !
186
        descending.push_back(false);
×
187
    } else if (PyBool_Check(descending_obj) || PyLong_Check(descending_obj)) {
80 ✔
188
        descending.push_back(PyObject_IsTrue(descending_obj) != 0);
62 !
189
    } else {
31 ✔
190
        PyObject* dseq = PySequence_Fast(
18 !
191
            descending_obj, "descending must be a bool or a sequence of bool");
9 ✔
192
        if (!dseq) return nullptr;
18 ✔
193
        const Py_ssize_t n = PySequence_Fast_GET_SIZE(dseq);
18 !
194
        for (Py_ssize_t i = 0; i < n; ++i)
70 ✔
195
            descending.push_back(
52 !
196
                PyObject_IsTrue(PySequence_Fast_GET_ITEM(dseq, i)) != 0);
52 !
197
        Py_DECREF(dseq);
9 !
198
    }
199
    return run_lazy_op([&] {
160 !
200
        return b->lf.sort_by_multi(std::move(names), std::move(descending));
80 !
201
    });
40 ✔
202
}
80 ✔
203

204
PyObject* LazyFrame_tail(PyObject* self, PyObject* n) {
2 ✔
205
    LazyFrameObject* b = as_lazyframe(self);
2 !
206
    if (!b) return nullptr;
2 ✔
207
    long long v = PyLong_AsLongLong(n);
2 !
208
    if (v == -1 && PyErr_Occurred()) return nullptr;
2 !
209
    return run_lazy_op([&] { return b->lf.tail(v); });
4 !
210
}
1 ✔
211

212
PyObject* LazyFrame_drop_nulls(PyObject* self, PyObject*) {
×
213
    LazyFrameObject* b = as_lazyframe(self);
×
214
    if (!b) return nullptr;
×
215
    return run_lazy_op([&] { return b->lf.drop_nulls(); });
×
216
}
217

218
PyObject* LazyFrame_fill_null(PyObject* self, PyObject* value) {
×
219
    LazyFrameObject* b = as_lazyframe(self);
×
220
    if (!b) return nullptr;
×
221
    dftu_scalar s{};
×
222
    if (!py_to_scalar(value, &s)) return nullptr;
×
223
    return run_lazy_op([&] { return b->lf.fill_null(s); });
×
224
}
225

226
PyObject* LazyFrame_with_row_index(PyObject* self, PyObject* name) {
30 ✔
227
    LazyFrameObject* b = as_lazyframe(self);
30 !
228
    if (!b) return nullptr;
30 ✔
229
    const char* s = PyUnicode_AsUTF8(name);
30 !
230
    if (!s) return nullptr;
30 ✔
231
    return run_lazy_op([&] { return b->lf.with_row_index(s); });
60 !
232
}
15 ✔
233

234
PyObject* LazyFrame_null_count(PyObject* self, PyObject*) {
×
235
    LazyFrameObject* b = as_lazyframe(self);
×
236
    if (!b) return nullptr;
×
237
    return run_lazy_op([&] { return b->lf.null_count(); });
×
238
}
239

240
PyObject* LazyFrame_reduce(PyObject* self, PyObject* arg) {
2 ✔
241
    LazyFrameObject* b = as_lazyframe(self);
2 !
242
    if (!b) return nullptr;
2 ✔
243
    const char* agg = PyUnicode_AsUTF8(arg);
2 !
244
    if (!agg) return nullptr;
2 ✔
245
    return run_lazy_op(
3 !
246
        [&] { return b->lf.reduce(dataframe::agg_from_string(agg)); });
4 !
247
}
1 ✔
248

249
PyObject* LazyFrame_group_transform(PyObject* self, PyObject* args,
32 ✔
250
                                    PyObject* kwds) {
251
    LazyFrameObject* b = as_lazyframe(self);
32 !
252
    if (!b) return nullptr;
32 ✔
253
    dftracer::utils::python::GroupwiseArgs a;
32 ✔
254
    if (!dftracer::utils::python::parse_groupwise_args(
32 !
255
            args, kwds, strings_from_str_or_seq, a))
16 ✔
256
        return nullptr;
×
257
    return run_lazy_op([&] {
64 !
258
        return b->lf.group_by(a.keys).transform(
48 !
259
            static_cast<dataframe::GroupwiseOp>(a.kind), a.n,
32 ✔
260
            static_cast<dataframe::RankMethod>(a.method), a.ascending != 0);
48 !
261
    });
16 ✔
262
}
32 ✔
263

264
// reduce_specs(agg, keys) -> list[str]: the "op:column:out" specs that
265
// broadcast `agg` over every eligible non-key column of the plan's schema.
266
PyObject* LazyFrame_reduce_specs(PyObject* self, PyObject* args) {
48 ✔
267
    LazyFrameObject* b = as_lazyframe(self);
48 !
268
    if (!b) return nullptr;
48 ✔
269
    const char* agg = nullptr;
48 ✔
270
    PyObject* keys_obj = Py_None;
48 ✔
271
    if (!PyArg_ParseTuple(args, "s|O", &agg, &keys_obj)) return nullptr;
48 !
272
    std::vector<std::string> keys;
48 ✔
273
    if (keys_obj != Py_None && !strings_from_str_or_seq(keys_obj, keys))
48 !
274
        return nullptr;
×
275
    std::vector<std::string> specs;
48 ✔
276
    try {
277
        specs = dftracer::utils::python::reduce_spec_strings(
72 !
278
            b->lf.reduce_specs(dataframe::agg_from_string(agg), keys));
96 !
279
    } catch (const std::out_of_range& e) {
24 !
280
        PyErr_SetString(PyExc_KeyError, e.what());
×
281
        return nullptr;
×
282
    } catch (const std::exception& e) {
×
283
        PyErr_SetString(PyExc_ValueError, e.what());
×
284
        return nullptr;
×
285
    }
×
286
    PyObject* list = PyList_New(static_cast<Py_ssize_t>(specs.size()));
48 !
287
    if (!list) return nullptr;
48 ✔
288
    for (std::size_t i = 0; i < specs.size(); ++i) {
142 ✔
289
        PyObject* s = PyUnicode_FromString(specs[i].c_str());
94 !
290
        if (!s) {
94 !
291
            Py_DECREF(list);
×
292
            return nullptr;
×
293
        }
294
        PyList_SET_ITEM(list, static_cast<Py_ssize_t>(i), s);
94 !
295
    }
47 ✔
296
    return list;
48 ✔
297
}
48 ✔
298

299
PyObject* LazyFrame_explode(PyObject* self, PyObject* column) {
×
300
    LazyFrameObject* b = as_lazyframe(self);
×
301
    if (!b) return nullptr;
×
302
    const char* s = PyUnicode_AsUTF8(column);
×
303
    if (!s) return nullptr;
×
304
    return run_lazy_op([&] { return b->lf.explode(s); });
×
305
}
306

307
PyObject* LazyFrame_unnest(PyObject* self, PyObject* args, PyObject* kwds) {
6 ✔
308
    LazyFrameObject* b = as_lazyframe(self);
6 !
309
    if (!b) return nullptr;
6 ✔
310
    const char* column = nullptr;
6 ✔
311
    int keep_empty = 0;
6 ✔
312
    static const char* kw[] = {"column", "keep_empty", nullptr};
313
    if (!PyArg_ParseTupleAndKeywords(args, kwds, "s|p", const_cast<char**>(kw),
6 !
314
                                     &column, &keep_empty))
315
        return nullptr;
×
316
    return run_lazy_op([&] { return b->lf.unnest(column, keep_empty != 0); });
13 !
317
}
3 ✔
318

319
PyObject* LazyFrame_compare_agg(PyObject* self, PyObject* args) {
2 ✔
320
    LazyFrameObject* b = as_lazyframe(self);
2 !
321
    if (!b) return nullptr;
2 ✔
322
    PyObject* other = nullptr;
2 ✔
323
    Py_ssize_t n_key = 1;
2 ✔
324
    if (!PyArg_ParseTuple(args, "On", &other, &n_key)) return nullptr;
2 !
325
    LazyFrameObject* o = as_lazyframe(other);
2 !
326
    if (!o) return nullptr;
2 ✔
327
    return run_lazy_op([&] {
4 !
328
        return b->lf.compare_agg(o->lf, static_cast<std::int64_t>(n_key));
2 !
329
    });
1 ✔
330
}
1 ✔
331

332
// The output names of a two-frame kernel that keeps every left column and
333
// appends the right's columns other than `dropped` (the join keys), suffixing
334
// a right name that collides with a left one by "_right": the asof_join and
335
// interval_join contract. Empty when either schema is unknown.
336
std::vector<std::string> two_frame_out_names(
16 ✔
337
    const LazyFrame& left, const LazyFrame& right,
338
    const std::vector<std::string>& dropped) {
339
    std::vector<std::string> names = left.schema();
16 !
340
    const std::vector<std::string> rnames = right.schema();
16 !
341
    if (names.empty() || rnames.empty()) return {};
16 !
342
    for (const std::string& r : rnames) {
64 ✔
343
        if (std::find(dropped.begin(), dropped.end(), r) != dropped.end())
48 !
344
            continue;
32 ✔
345
        const bool collides =
8 ✔
346
            std::find(names.begin(), names.end(), r) != names.end();
16 !
347
        names.push_back(collides ? r + "_right" : r);
16 !
348
    }
349
    return names;
16 ✔
350
}
16 ✔
351

352
// The four Arrow-layer relational kernels as plan steps: each builds the C
353
// operand bag its dftu.frame.* row takes and appends it through frame_op,
354
// which copies every operand into the plan.
355
PyObject* LazyFrame_window(PyObject* self, PyObject* args, PyObject* kwds) {
108 ✔
356
    LazyFrameObject* b = as_lazyframe(self);
108 !
357
    if (!b) return nullptr;
108 !
358
    static const char* kw[] = {"partition_by", "order_by", "specs", nullptr};
359
    PyObject* part = nullptr;
108 ✔
360
    PyObject* order = nullptr;
108 ✔
361
    PyObject* specs = nullptr;
108 ✔
362
    if (!PyArg_ParseTupleAndKeywords(args, kwds, "OOO", const_cast<char**>(kw),
108 !
363
                                     &part, &order, &specs))
364
        return nullptr;
×
365
    std::vector<std::string> pcols;
108 ✔
366
    std::vector<std::string> ocols;
108 ✔
367
    if (!parse_string_seq(part, "window: partition_by must be names", pcols) ||
216 !
368
        !parse_string_seq(order, "window: order_by must be names", ocols))
108 !
369
        return nullptr;
×
370
    dftracer::utils::python::WindowSpecs parsed;
108 !
371
    if (!dftracer::utils::python::parse_window_specs(specs, parsed))
108 !
372
        return nullptr;
×
373
    std::vector<const char*> pc;
108 ✔
374
    std::vector<const char*> oc;
108 ✔
375
    for (const std::string& s : pcols) pc.push_back(s.c_str());
216 ✔
376
    for (const std::string& s : ocols) oc.push_back(s.c_str());
216 ✔
377
    dataframe::OpArgs a;
108 ✔
378
    a.strlist(1, pc.data(), static_cast<std::int32_t>(pc.size()))
108 ✔
379
        .strlist(2, oc.data(), static_cast<std::int32_t>(oc.size()))
108 ✔
380
        .winlist(3, parsed.specs);
108 !
381
    return run_lazy_op([&] {
162 !
382
        // A window keeps every input column and appends one per spec.
383
        std::vector<std::string> names = b->lf.schema();
108 !
384
        if (!names.empty())
108 ✔
385
            for (const dftu_window_spec& s : parsed.specs)
326 ✔
386
                names.emplace_back(s.out);
272 !
387
        return b->lf.frame_op("dftu.frame.window", a, {}, std::move(names));
162 !
388
    });
162 ✔
389
}
108 ✔
390

391
// column_op(column, op, column2=None, a=0, b=0, text=None): the registered
392
// column op `op` over `column` as a plan step (dftu.frame.column_op, a
393
// breaker), the column replaced in place.
394
PyObject* LazyFrame_column_op(PyObject* self, PyObject* args, PyObject* kwds) {
×
395
    LazyFrameObject* b = as_lazyframe(self);
×
396
    if (!b) return nullptr;
×
397
    static const char* kw[] = {"column", "op",   "column2", "a",
398
                               "b",      "text", nullptr};
399
    const char* column = nullptr;
×
400
    const char* op = nullptr;
×
401
    const char* column2 = nullptr;
×
402
    PyObject* a_obj = nullptr;
×
403
    PyObject* b_obj = nullptr;
×
404
    const char* text = nullptr;
×
405
    if (!PyArg_ParseTupleAndKeywords(args, kwds, "ss|zOOz",
×
406
                                     const_cast<char**>(kw), &column, &op,
407
                                     &column2, &a_obj, &b_obj, &text))
408
        return nullptr;
×
409
    dftu_scalar sa{};
×
410
    dftu_scalar sb{};
×
411
    sa.kind = DFTU_SCALAR_TAG_I64;
×
412
    sb.kind = DFTU_SCALAR_TAG_I64;
×
413
    if (a_obj && a_obj != Py_None && !py_to_scalar(a_obj, &sa)) return nullptr;
×
414
    if (b_obj && b_obj != Py_None && !py_to_scalar(b_obj, &sb)) return nullptr;
×
415
    const std::string c2 = column2 ? column2 : "";
×
416
    const std::string tx = text ? text : "";
×
417
    dataframe::OpArgs a;
×
418
    a.str(1, column).str(2, op).str(3, c2).scalar(4, sa).scalar(5, sb).str(6,
×
419
                                                                           tx);
420
    return run_lazy_op([&] {
×
421
        return b->lf.frame_op("dftu.frame.column_op", a, {}, b->lf.schema());
×
422
    });
×
423
}
×
424

425
PyObject* LazyFrame_gap_fill(PyObject* self, PyObject* args, PyObject* kwds) {
8 ✔
426
    LazyFrameObject* b = as_lazyframe(self);
8 !
427
    if (!b) return nullptr;
8 ✔
428
    static const char* kw[] = {"partition_by", "time",  "bucket", "values",
429
                               "mode",         "start", "end",    nullptr};
430
    PyObject* part = nullptr;
8 ✔
431
    const char* time = nullptr;
8 ✔
432
    long long bucket = 0;
8 ✔
433
    PyObject* values = nullptr;
8 ✔
434
    const char* mode = nullptr;
8 ✔
435
    PyObject* start_obj = Py_None;
8 ✔
436
    PyObject* end_obj = Py_None;
8 ✔
437
    if (!PyArg_ParseTupleAndKeywords(
8 !
438
            args, kwds, "OsLOs|OO", const_cast<char**>(kw), &part, &time,
4 ✔
439
            &bucket, &values, &mode, &start_obj, &end_obj))
440
        return nullptr;
×
441
    dftu_gap_fill_mode m;
442
    if (!dftracer::utils::python::gap_fill_mode_from_str(mode, &m))
8 !
443
        return nullptr;
×
444
    std::vector<std::string> pcols;
8 ✔
445
    std::vector<std::string> vcols;
8 ✔
446
    if (!parse_string_seq(part, "gap_fill: partition_by must be names",
8 !
447
                          pcols) ||
16 !
448
        !parse_string_seq(values, "gap_fill: values must be names", vcols))
8 !
449
        return nullptr;
×
450
    std::vector<std::int64_t> range;
8 ✔
451
    if (start_obj != Py_None) {
8 ✔
452
        const long long s = PyLong_AsLongLong(start_obj);
2 !
453
        const long long e = PyLong_AsLongLong(end_obj);
2 !
454
        if (PyErr_Occurred()) return nullptr;
2 !
455
        range = {static_cast<std::int64_t>(s), static_cast<std::int64_t>(e)};
2 !
456
    }
1 ✔
457
    std::vector<const char*> pc;
8 ✔
458
    std::vector<const char*> vc;
8 ✔
459
    for (const std::string& s : pcols) pc.push_back(s.c_str());
16 ✔
460
    for (const std::string& s : vcols) vc.push_back(s.c_str());
16 ✔
461
    dataframe::OpArgs a;
8 ✔
462
    a.strlist(1, pc.data(), static_cast<std::int32_t>(pc.size()))
8 ✔
463
        .str(2, time)
8 !
464
        .i64(3, static_cast<std::int64_t>(bucket))
8 ✔
465
        .strlist(4, vc.data(), static_cast<std::int32_t>(vc.size()))
8 ✔
466
        .i32(5, static_cast<std::int32_t>(m))
8 ✔
467
        .i64list(6, range);
8 !
468
    return run_lazy_op([&] {
12 !
469
        return b->lf.frame_op("dftu.frame.gap_fill", a, {}, b->lf.schema());
8 !
470
    });
4 ✔
471
}
8 ✔
472

473
PyObject* LazyFrame_asof(PyObject* self, PyObject* args, PyObject* kwds) {
14 ✔
474
    LazyFrameObject* b = as_lazyframe(self);
14 !
475
    if (!b) return nullptr;
14 !
476
    static const char* kw[] = {"other",     "on",        "by",
477
                               "direction", "tolerance", nullptr};
478
    PyObject* other = nullptr;
14 ✔
479
    const char* on = nullptr;
14 ✔
480
    PyObject* by = nullptr;
14 ✔
481
    const char* direction = nullptr;
14 ✔
482
    PyObject* tol_obj = Py_None;
14 ✔
483
    if (!PyArg_ParseTupleAndKeywords(args, kwds, "OsOsO",
14 !
484
                                     const_cast<char**>(kw), &other, &on, &by,
485
                                     &direction, &tol_obj))
486
        return nullptr;
×
487
    LazyFrameObject* o = as_lazyframe(other);
14 !
488
    if (!o) return nullptr;
14 ✔
489
    dftu_asof_direction dir;
490
    if (!dftracer::utils::python::asof_direction_from_str(direction, &dir))
14 ✔
491
        return nullptr;
2 ✔
492
    std::vector<std::string> equi;
12 ✔
493
    if (!parse_string_seq(by, "asof: by must be names", equi)) return nullptr;
12 !
494
    std::int64_t tol = -1;
12 ✔
495
    if (tol_obj != Py_None) {
12 ✔
496
        const long long t = PyLong_AsLongLong(tol_obj);
2 !
497
        if (PyErr_Occurred()) return nullptr;
2 !
498
        if (t < 0) {
2 !
499
            PyErr_SetString(PyExc_ValueError,
×
500
                            "asof: tolerance must not be negative");
501
            return nullptr;
×
502
        }
503
        tol = static_cast<std::int64_t>(t);
2 ✔
504
    }
1 ✔
505
    std::vector<const char*> ec;
12 ✔
506
    for (const std::string& s : equi) ec.push_back(s.c_str());
24 ✔
507
    dataframe::OpArgs a;
12 ✔
508
    a.str(2, on)
12 !
509
        .strlist(3, ec.data(), static_cast<std::int32_t>(ec.size()))
12 ✔
510
        .i32(4, static_cast<std::int32_t>(dir))
12 ✔
511
        .i64(5, tol);
12 ✔
512
    return run_lazy_op([&] {
18 !
513
        std::vector<std::string> dropped = equi;
12 !
514
        dropped.emplace_back(on);
12 !
515
        return b->lf.frame_op("dftu.frame.asof", a, {o->lf},
18 !
516
                              two_frame_out_names(b->lf, o->lf, dropped));
24 !
517
    });
18 ✔
518
}
13 ✔
519

520
PyObject* LazyFrame_interval(PyObject* self, PyObject* args, PyObject* kwds) {
4 ✔
521
    LazyFrameObject* b = as_lazyframe(self);
4 !
522
    if (!b) return nullptr;
4 ✔
523
    static const char* kw[] = {"other", "point", "lo",   "hi",
524
                               "by",    "outer", nullptr};
525
    PyObject* other = nullptr;
4 ✔
526
    const char* point = nullptr;
4 ✔
527
    const char* lo = nullptr;
4 ✔
528
    const char* hi = nullptr;
4 ✔
529
    PyObject* by = nullptr;
4 ✔
530
    int outer = 0;
4 ✔
531
    if (!PyArg_ParseTupleAndKeywords(args, kwds, "OsssOp",
4 !
532
                                     const_cast<char**>(kw), &other, &point,
533
                                     &lo, &hi, &by, &outer))
534
        return nullptr;
×
535
    LazyFrameObject* o = as_lazyframe(other);
4 !
536
    if (!o) return nullptr;
4 ✔
537
    std::vector<std::string> equi;
4 ✔
538
    if (!parse_string_seq(by, "interval: by must be names", equi))
4 !
539
        return nullptr;
×
540
    std::vector<const char*> ec;
4 ✔
541
    for (const std::string& s : equi) ec.push_back(s.c_str());
4 !
542
    dataframe::OpArgs a;
4 ✔
543
    a.str(2, point)
4 !
544
        .str(3, lo)
4 !
545
        .str(4, hi)
4 !
546
        .strlist(5, ec.data(), static_cast<std::int32_t>(ec.size()))
4 ✔
547
        .i32(6, outer != 0 ? 1 : 0);
4 ✔
548
    return run_lazy_op([&] {
6 !
549
        std::vector<std::string> dropped = equi;
4 !
550
        dropped.emplace_back(lo);
4 !
551
        dropped.emplace_back(hi);
4 !
552
        return b->lf.frame_op("dftu.frame.interval", a, {o->lf},
6 !
553
                              two_frame_out_names(b->lf, o->lf, dropped));
8 !
554
    });
6 ✔
555
}
4 ✔
556

557
// unpivot(id_vars, value_vars); melt is the same call.
558
PyObject* LazyFrame_unpivot(PyObject* self, PyObject* args) {
×
559
    LazyFrameObject* b = as_lazyframe(self);
×
560
    if (!b) return nullptr;
×
561
    PyObject* id_obj = nullptr;
×
562
    PyObject* val_obj = nullptr;
×
563
    if (!PyArg_ParseTuple(args, "OO", &id_obj, &val_obj)) return nullptr;
×
564
    std::vector<std::string> id_vars, value_vars;
×
565
    if (!parse_str_list(id_obj, "id_vars", id_vars) ||
×
566
        !parse_str_list(val_obj, "value_vars", value_vars))
×
567
        return nullptr;
×
568
    return run_lazy_op([&] {
×
569
        return b->lf.unpivot(std::move(id_vars), std::move(value_vars));
×
570
    });
×
571
}
×
572

573
PyObject* LazyFrame_topk(PyObject* self, PyObject* args, PyObject* kwds) {
8 ✔
574
    LazyFrameObject* b = as_lazyframe(self);
8 !
575
    if (!b) return nullptr;
8 ✔
576
    const char* name = nullptr;
8 ✔
577
    long long k = 0;
8 ✔
578
    int largest = 1;
8 ✔
579
    static const char* kw[] = {"name", "k", "largest", nullptr};
580
    if (!PyArg_ParseTupleAndKeywords(args, kwds, "sL|p", const_cast<char**>(kw),
8 !
581
                                     &name, &k, &largest))
582
        return nullptr;
×
583
    return run_lazy_op([&] { return b->lf.topk(name, k, largest != 0); });
16 !
584
}
4 ✔
585

586
// group_by(key, aggs): key is a column name or a sequence of names (composite
587
// key); aggs is a sequence of "op[:column]" string specs.
588
PyObject* LazyFrame_group_by(PyObject* self, PyObject* args) {
176 ✔
589
    LazyFrameObject* b = as_lazyframe(self);
176 !
590
    if (!b) return nullptr;
176 ✔
591
    PyObject* key_obj = nullptr;
176 ✔
592
    PyObject* aggs_obj = nullptr;
176 ✔
593
    if (!PyArg_ParseTuple(args, "OO", &key_obj, &aggs_obj)) return nullptr;
176 !
594
    std::vector<std::string> keys;
176 ✔
595
    if (!strings_from_str_or_seq(key_obj, keys)) return nullptr;
176 !
596
    std::vector<GroupAgg> aggs;
176 ✔
597
    if (!aggs_from_seq(aggs_obj, aggs)) return nullptr;
176 !
598
    return run_lazy_op(
264 !
599
        [&] { return b->lf.group_by(std::move(keys), std::move(aggs)); });
352 !
600
}
176 ✔
601

602
PyObject* LazyFrame_sort_by(PyObject* self, PyObject* args, PyObject* kwds) {
56 ✔
603
    LazyFrameObject* b = as_lazyframe(self);
56 !
604
    if (!b) return nullptr;
56 ✔
605
    const char* name = nullptr;
56 ✔
606
    int descending = 0;
56 ✔
607
    static const char* kw[] = {"name", "descending", nullptr};
608
    if (!PyArg_ParseTupleAndKeywords(args, kwds, "s|p", const_cast<char**>(kw),
56 !
609
                                     &name, &descending))
610
        return nullptr;
×
611
    return run_lazy_op([&] { return b->lf.sort_by(name, descending != 0); });
112 !
612
}
28 ✔
613

614
PyObject* LazyFrame_join(PyObject* self, PyObject* args, PyObject* kwds) {
154 ✔
615
    LazyFrameObject* b = as_lazyframe(self);
154 !
616
    if (!b) return nullptr;
154 ✔
617
    PyObject* other = nullptr;
154 ✔
618
    PyObject* on = nullptr;
154 ✔
619
    const char* how = "inner";
154 ✔
620
    PyObject* left_on = nullptr;
154 ✔
621
    PyObject* right_on = nullptr;
154 ✔
622
    const char* suffix = "_right";
154 ✔
623
    static const char* kw[] = {"other",    "on",     "how",  "left_on",
624
                               "right_on", "suffix", nullptr};
625
    if (!PyArg_ParseTupleAndKeywords(args, kwds, "O|OsOOs",
154 !
626
                                     const_cast<char**>(kw), &other, &on, &how,
627
                                     &left_on, &right_on, &suffix))
628
        return nullptr;
×
629
    LazyFrameObject* o = as_lazyframe(other);
154 !
630
    if (!o) return nullptr;
154 ✔
631
    dataframe::JoinHow jh;
632
    if (!dftracer::utils::python::join_how_from_str(how, jh)) return nullptr;
154 ✔
633
    std::vector<std::string> l;
152 ✔
634
    std::vector<std::string> r;
152 ✔
635
    if (!dftracer::utils::python::join_keys_from_objs(on, left_on, right_on, jh,
152 !
636
                                                      l, r))
637
        return nullptr;
×
638
    return run_lazy_op([&] {
304 !
639
        return b->lf.join(o->lf, std::move(l), std::move(r), jh, suffix);
152 !
640
    });
76 ✔
641
}
153 ✔
642

643
PyObject* LazyFrame_concat(PyObject* self, PyObject* args) {
16 ✔
644
    LazyFrameObject* b = as_lazyframe(self);
16 !
645
    if (!b) return nullptr;
16 ✔
646
    const Py_ssize_t n = PyTuple_GET_SIZE(args);
16 ✔
647
    std::vector<LazyFrameObject*> others;
16 ✔
648
    others.reserve(static_cast<std::size_t>(n));
16 !
649
    for (Py_ssize_t i = 0; i < n; ++i) {
32 ✔
650
        LazyFrameObject* o = as_lazyframe(PyTuple_GET_ITEM(args, i));
18 !
651
        if (!o) return nullptr;
18 ✔
652
        others.push_back(o);
16 !
653
    }
8 ✔
654
    return run_lazy_op([&] {
28 !
655
        dataframe::LazyFrame out = b->lf;
14 ✔
656
        for (LazyFrameObject* o : others) out = out.concat(o->lf);
29 ✔
657
        return out;
12 ✔
658
    });
15 !
659
}
16 ✔
660

661
PyObject* LazyFrame_unique(PyObject* self, PyObject* args, PyObject* kwds) {
70 ✔
662
    LazyFrameObject* b = as_lazyframe(self);
70 !
663
    if (!b) return nullptr;
70 ✔
664
    PyObject* subset_obj = Py_None;
70 ✔
665
    static const char* kw[] = {"subset", nullptr};
666
    if (!PyArg_ParseTupleAndKeywords(args, kwds, "|O", const_cast<char**>(kw),
70 !
667
                                     &subset_obj))
668
        return nullptr;
×
669
    std::vector<std::string> subset;
70 ✔
670
    if (subset_obj != Py_None && !strings_from_str_or_seq(subset_obj, subset))
70 !
671
        return nullptr;
×
672
    return run_lazy_op([&] { return b->lf.unique(std::move(subset)); });
140 !
673
}
70 ✔
674

675
PyObject* LazyFrame_sample(PyObject* self, PyObject* args, PyObject* kwds) {
×
676
    LazyFrameObject* b = as_lazyframe(self);
×
677
    if (!b) return nullptr;
×
678
    long long n = 0;
×
679
    unsigned long long seed = 0;
×
680
    static const char* kw[] = {"n", "seed", nullptr};
681
    if (!PyArg_ParseTupleAndKeywords(args, kwds, "L|K", const_cast<char**>(kw),
×
682
                                     &n, &seed))
683
        return nullptr;
×
684
    return run_lazy_op([&] { return b->lf.sample(n, seed); });
×
685
}
686

687
PyObject* LazyFrame_is_duplicated(PyObject* self, PyObject*) {
×
688
    LazyFrameObject* b = as_lazyframe(self);
×
689
    if (!b) return nullptr;
×
690
    return run_lazy_op([&] { return b->lf.is_duplicated(); });
×
691
}
692

693
PyObject* LazyFrame_is_unique(PyObject* self, PyObject*) {
×
694
    LazyFrameObject* b = as_lazyframe(self);
×
695
    if (!b) return nullptr;
×
696
    return run_lazy_op([&] { return b->lf.is_unique(); });
×
697
}
698

699
// group_by_dynamic(time_col, every, period=None, aggs=[], origin=0,
700
// origin_min=False). period defaults to every (tumbling window).
701
PyObject* LazyFrame_group_by_dynamic(PyObject* self, PyObject* args,
4 ✔
702
                                     PyObject* kwds) {
703
    LazyFrameObject* b = as_lazyframe(self);
4 !
704
    if (!b) return nullptr;
4 ✔
705
    const char* time_col = nullptr;
4 ✔
706
    long long every = 0;
4 ✔
707
    PyObject* period_obj = Py_None;
4 ✔
708
    PyObject* aggs_obj = nullptr;
4 ✔
709
    long long origin = 0;
4 ✔
710
    int origin_min = 0;
4 ✔
711
    static const char* kw[] = {"time_col", "every",      "period", "aggs",
712
                               "origin",   "origin_min", nullptr};
713
    if (!PyArg_ParseTupleAndKeywords(
4 !
714
            args, kwds, "sL|OOLp", const_cast<char**>(kw), &time_col, &every,
2 ✔
715
            &period_obj, &aggs_obj, &origin, &origin_min))
716
        return nullptr;
×
717
    long long period = every;
4 ✔
718
    if (period_obj && period_obj != Py_None) {
4 !
719
        period = PyLong_AsLongLong(period_obj);
×
720
        if (period == -1 && PyErr_Occurred()) return nullptr;
×
721
    }
722
    std::vector<GroupAgg> aggs;
4 ✔
723
    if (aggs_obj && aggs_obj != Py_None && !aggs_from_seq(aggs_obj, aggs))
4 !
724
        return nullptr;
×
725
    return run_lazy_op([&] {
8 !
726
        return b->lf.group_by_dynamic(time_col, every, period, std::move(aggs),
10 !
727
                                      origin, origin_min != 0);
4 !
728
    });
2 ✔
729
}
4 ✔
730

731
PyObject* LazyFrame_pivot(PyObject* self, PyObject* args, PyObject* kwds) {
×
732
    LazyFrameObject* b = as_lazyframe(self);
×
733
    if (!b) return nullptr;
×
734
    const char* index = nullptr;
×
735
    const char* on = nullptr;
×
736
    const char* values = nullptr;
×
737
    const char* agg = "first";
×
738
    static const char* kw[] = {"index", "on", "values", "agg", nullptr};
739
    if (!PyArg_ParseTupleAndKeywords(args, kwds, "sss|s",
×
740
                                     const_cast<char**>(kw), &index, &on,
741
                                     &values, &agg))
742
        return nullptr;
×
743
    return run_lazy_op([&] { return b->lf.pivot(index, on, values, agg); });
×
744
}
745

746
PyObject* LazyFrame_to_dummies(PyObject* self, PyObject* column) {
2 ✔
747
    LazyFrameObject* b = as_lazyframe(self);
2 !
748
    if (!b) return nullptr;
2 ✔
749
    const char* s = PyUnicode_AsUTF8(column);
2 !
750
    if (!s) return nullptr;
2 ✔
751
    return run_lazy_op([&] { return b->lf.to_dummies(s); });
4 !
752
}
1 ✔
753

754
PyObject* LazyFrame_describe(PyObject* self, PyObject*) {
×
755
    LazyFrameObject* b = as_lazyframe(self);
×
756
    if (!b) return nullptr;
×
757
    return run_lazy_op([&] { return b->lf.describe(); });
×
758
}
759

760
PyObject* LazyFrame_memory_budget(PyObject* self, PyObject* bytes) {
2 ✔
761
    LazyFrameObject* b = as_lazyframe(self);
2 !
762
    if (!b) return nullptr;
2 ✔
763
    unsigned long long v = PyLong_AsUnsignedLongLong(bytes);
2 !
764
    if (v == static_cast<unsigned long long>(-1) && PyErr_Occurred())
2 !
765
        return nullptr;
×
766
    return run_lazy_op([&] { return b->lf.memory_budget(v); });
4 !
767
}
1 ✔
768

769
PyObject* LazyFrame_auto_spill(PyObject* self, PyObject*) {
×
770
    LazyFrameObject* b = as_lazyframe(self);
×
771
    if (!b) return nullptr;
×
772
    return run_lazy_op([&] { return b->lf.auto_spill(); });
×
773
}
774

775
PyObject* LazyFrame_schema(PyObject* self, PyObject*) {
1,450 ✔
776
    LazyFrameObject* b = as_lazyframe(self);
1,450 !
777
    if (!b) return nullptr;
1,450 ✔
778
    std::vector<std::string> names;
1,450 ✔
779
    try {
780
        names = b->lf.schema();
1,450 ✔
781
    } catch (const std::exception& e) {
727 !
782
        PyErr_SetString(PyExc_ValueError, e.what());
4 !
783
        return nullptr;
4 ✔
784
    }
4 !
785
    return str_list_from(names);
1,446 !
786
}
1,452 ✔
787

788
PyObject* LazyFrame_explain(PyObject* self, PyObject*) {
20 ✔
789
    LazyFrameObject* b = as_lazyframe(self);
20 !
790
    if (!b) return nullptr;
20 ✔
791
    std::string text;
20 ✔
792
    try {
793
        text = b->lf.explain();
20 !
794
    } catch (const std::exception& e) {
10 !
795
        PyErr_SetString(PyExc_ValueError, e.what());
×
796
        return nullptr;
×
797
    }
×
798
    return PyUnicode_FromStringAndSize(text.data(),
30 !
799
                                       static_cast<Py_ssize_t>(text.size()));
30 ✔
800
}
20 ✔
801

802
PyObject* LazyFrame_collect(PyObject* self, PyObject* args, PyObject* kwds) {
646 ✔
803
    LazyFrameObject* b = as_lazyframe(self);
646 !
804
    if (!b) return nullptr;
646 ✔
805
    long long morsel_rows = 0;
646 ✔
806
    PyObject* runtime_arg = nullptr;
646 ✔
807
    static const char* kw[] = {"morsel_rows", "runtime", nullptr};
808
    if (!PyArg_ParseTupleAndKeywords(args, kwds, "|LO", const_cast<char**>(kw),
646 !
809
                                     &morsel_rows, &runtime_arg))
810
        return nullptr;
×
811
    std::shared_ptr<dftracer::utils::Runtime> rt =
812
        runtime_from_arg(runtime_arg);
646 !
813
    if (!rt) return nullptr;
646 ✔
814
    dataframe::DataFrame out;
646 ✔
815
    if (!run_blocking(
646 !
816
            [&] { out = rt->submit(b->lf.collect(morsel_rows)).get(); }))
969 !
817
        return nullptr;
×
818
    return dftracer::utils::python::wrap_dataframe(std::move(out));
646 !
819
}
646 ✔
820

821
// Drains the plan's chunk generator into `state`, which the Python iterator
822
// pulls from with the GIL released.
823
dftracer::utils::coro::CoroTask<void> drain_stream(
278 !
824
    dftracer::utils::CoroScope&,
825
    std::shared_ptr<
826
        dftracer::utils::python::StreamingState<dataframe::DataFrame>>
827
        state,
828
    LazyFrame lf, std::int64_t morsel_rows) {
15 !
829
    try {
830
        auto gen = lf.stream(morsel_rows);
15 !
831
        while (auto df = co_await gen.next()) {
203 !
832
            if (state->cancelled()) break;
32 !
833
            const std::size_t bytes =
64 ✔
834
                static_cast<std::size_t>(df->num_rows()) *
64 ✔
835
                static_cast<std::size_t>(df->num_columns() + 1) * 16;
64 ✔
836
            if (!state->push(std::move(*df), bytes)) break;
32 !
837
        }
47 !
838
        state->complete();
15 !
839
    } catch (...) {
109 ✔
840
        state->fail(std::current_exception());
×
841
    }
×
842
}
218 !
843

844
PyObject* LazyFrame_stream(PyObject* self, PyObject* args, PyObject* kwds) {
30 ✔
845
    LazyFrameObject* b = as_lazyframe(self);
30 !
846
    if (!b) return nullptr;
30 ✔
847
    long long morsel_rows = 0;
30 ✔
848
    PyObject* runtime_arg = nullptr;
30 ✔
849
    static const char* kw[] = {"morsel_rows", "runtime", nullptr};
850
    if (!PyArg_ParseTupleAndKeywords(args, kwds, "|LO", const_cast<char**>(kw),
30 !
851
                                     &morsel_rows, &runtime_arg))
NEW
852
        return nullptr;
×
853
    std::shared_ptr<dftracer::utils::Runtime> rt =
854
        runtime_from_arg(runtime_arg);
30 !
855
    if (!rt) return nullptr;
30 ✔
856
    namespace py = dftracer::utils::python;
857
    auto state = std::make_shared<py::StreamingState<dataframe::DataFrame>>(
15 !
858
        dftracer::utils::compute_memory_budget(0));
30 !
859
    auto* it = reinterpret_cast<py::ArrowStreamingIteratorObject*>(
15 ✔
860
        py::ArrowStreamingIteratorType.tp_new(&py::ArrowStreamingIteratorType,
30 !
861
                                              nullptr, nullptr));
862
    if (!it) return nullptr;
30 ✔
863
    it->cpp_state->state = state;
30 ✔
864
    it->cpp_state->pull_df = [state]() { return state->pull(); };
124 !
865
    it->cpp_state->get_error = [state]() { return state->error(); };
60 !
866
    it->cpp_state->cancel = [state]() { state->cancel(); };
60 !
867
    LazyFrame lf = b->lf;
30 !
868
    Py_BEGIN_ALLOW_THREADS rt->submit(
75 !
869
        dftracer::utils::run_coro_scope(rt->executor(), drain_stream, state,
45 !
870
                                        std::move(lf),
30 ✔
871
                                        static_cast<std::int64_t>(morsel_rows)),
15 ✔
872
        "lazyframe_stream");
15 !
873
    Py_END_ALLOW_THREADS return reinterpret_cast<PyObject*>(it);
30 !
874
}
30 ✔
875

876
PyObject* LazyFrame_output_schema(PyObject* self, PyObject*) {
4 ✔
877
    LazyFrameObject* b = as_lazyframe(self);
4 !
878
    if (!b) return nullptr;
4 ✔
879
    dataframe::Schema schema;
4 ✔
880
    try {
881
        schema = b->lf.output_schema();
4 !
882
    } catch (const std::exception& e) {
2 !
NEW
883
        PyErr_SetString(PyExc_ValueError, e.what());
×
NEW
884
        return nullptr;
×
NEW
885
    }
×
886
    PyObject* out = PyList_New(static_cast<Py_ssize_t>(schema.fields.size()));
4 !
887
    if (!out) return nullptr;
4 ✔
888
    for (std::size_t i = 0; i < schema.fields.size(); ++i) {
12 ✔
889
        const dataframe::Field& f = schema.fields[i];
8 ✔
890
        PyObject* item =
4 ✔
891
            Py_BuildValue("(si)", f.name.c_str(), static_cast<int>(f.type.id));
8 !
892
        if (!item) {
8 !
893
            Py_DECREF(out);
×
NEW
894
            return nullptr;
×
895
        }
896
        PyList_SET_ITEM(out, static_cast<Py_ssize_t>(i), item);
8 !
897
    }
4 ✔
898
    return out;
4 ✔
899
}
4 ✔
900

901
PyMethodDef LazyFrame_methods[] = {
902
    {"filter", LazyFrame_filter, METH_O,
903
     "filter(ast) -> LazyFrame keeping rows where the predicate holds."},
904
    {"with_column", LazyFrame_with_column, METH_VARARGS,
905
     "with_column(name, ast) -> LazyFrame with a column added or replaced."},
906
    {"select", LazyFrame_select, METH_O,
907
     "select(names) -> LazyFrame projected to the given columns."},
908
    {"rename", LazyFrame_rename, METH_O,
909
     "rename(names) -> LazyFrame with columns renamed positionally."},
910
    {"slice", LazyFrame_slice, METH_VARARGS,
911
     "slice(offset, len) -> LazyFrame of the row window."},
912
    {"head", LazyFrame_head, METH_O,
913
     "head(n) -> LazyFrame of the first n rows."},
914
    {"tail", LazyFrame_tail, METH_O,
915
     "tail(n) -> LazyFrame of the last n rows."},
916
    {"drop_nulls", LazyFrame_drop_nulls, METH_NOARGS,
917
     "drop_nulls() -> LazyFrame with rows holding any null removed."},
918
    {"fill_null", LazyFrame_fill_null, METH_O,
919
     "fill_null(value) -> LazyFrame with nulls filled in every column."},
920
    {"with_row_index", LazyFrame_with_row_index, METH_O,
921
     "with_row_index(name) -> LazyFrame with a prepended Int64 index column."},
922
    {"null_count", LazyFrame_null_count, METH_NOARGS,
923
     "null_count() -> LazyFrame of each column's null count (one row)."},
924
    {"reduce", LazyFrame_reduce, METH_O,
925
     "reduce(agg) -> LazyFrame: the aggregate named `agg` over every eligible "
926
     "column, a streaming one-group group-by."},
927
    {"group_transform", DFTU_PYCFUNCTION(LazyFrame_group_transform),
928
     METH_VARARGS | METH_KEYWORDS,
929
     "group_transform(keys, kind, n=0, method='average', ascending=True) -> "
930
     "LazyFrame: the group-wise transform `kind` (cumsum, cummax, cummin, "
931
     "cumcount, shift, diff, pct_change, rank, ngroup, head, tail, nth) over "
932
     "the groups of `keys`, one value per input row in input order."},
933
    {"reduce_specs", LazyFrame_reduce_specs, METH_VARARGS,
934
     "reduce_specs(agg, keys=None) -> list[str]: the 'op:column:out' specs "
935
     "that broadcast `agg` over every eligible non-key column, for group_by."},
936
    {"explode", LazyFrame_explode, METH_O,
937
     "explode(column) -> LazyFrame with the List column expanded per element."},
938
    {"unnest", DFTU_PYCFUNCTION(LazyFrame_unnest), METH_VARARGS | METH_KEYWORDS,
939
     "unnest(column, keep_empty=False) -> LazyFrame expanding a List column "
940
     "one row per element; an empty/null list drops the row unless "
941
     "keep_empty, and a List<Struct> flattens into one column per field."},
942
    {"compare_agg", LazyFrame_compare_agg, METH_VARARGS,
943
     "compare_agg(variant, n_key) -> LazyFrame: DataFrame.compare_agg over "
944
     "the two collected plans."},
945
    {"window", DFTU_PYCFUNCTION(LazyFrame_window), METH_VARARGS | METH_KEYWORDS,
946
     "window(partition_by, order_by, specs) -> LazyFrame: SQL window "
947
     "functions over the collected plan; specs are normalized 9-tuples."},
948
    {"column_op", DFTU_PYCFUNCTION(LazyFrame_column_op),
949
     METH_VARARGS | METH_KEYWORDS,
950
     "column_op(column, op, column2=None, a=0, b=0, text=None) -> the "
951
     "registered column op over `column` as a plan step (a breaker)."},
952
    {"gap_fill", DFTU_PYCFUNCTION(LazyFrame_gap_fill),
953
     METH_VARARGS | METH_KEYWORDS,
954
     "gap_fill(partition_by, time, bucket, values, mode, start=None, "
955
     "end=None) -> LazyFrame: a regular time grid over the collected plan."},
956
    {"asof", DFTU_PYCFUNCTION(LazyFrame_asof), METH_VARARGS | METH_KEYWORDS,
957
     "asof(other, on, by, direction, tolerance) -> LazyFrame: temporal "
958
     "nearest-match join of the two collected plans."},
959
    {"interval", DFTU_PYCFUNCTION(LazyFrame_interval),
960
     METH_VARARGS | METH_KEYWORDS,
961
     "interval(other, point, lo, hi, by, outer) -> LazyFrame: point-in-range "
962
     "join of the two collected plans."},
963
    {"unpivot", LazyFrame_unpivot, METH_VARARGS,
964
     "unpivot(id_vars, value_vars) -> LazyFrame reshaped wide to long."},
965
    {"melt", LazyFrame_unpivot, METH_VARARGS, "melt(...) -> alias of unpivot."},
966
    {"topk", DFTU_PYCFUNCTION(LazyFrame_topk), METH_VARARGS | METH_KEYWORDS,
967
     "topk(name, k, largest=True) -> LazyFrame of the k extreme rows."},
968
    {"group_by", LazyFrame_group_by, METH_VARARGS,
969
     "group_by(key, aggs) -> LazyFrame grouped by key with string-spec aggs."},
970
    {"sort_by", DFTU_PYCFUNCTION(LazyFrame_sort_by),
971
     METH_VARARGS | METH_KEYWORDS,
972
     "sort_by(name, descending=False) -> LazyFrame sorted by one column."},
973
    {"take", LazyFrame_take, METH_O,
974
     "take(indices) -> LazyFrame keeping the rows at indices (any order, "
975
     "repeats allowed)."},
976
    {"filter_mask", LazyFrame_filter_mask, METH_O,
977
     "filter_mask(mask) -> LazyFrame keeping rows where the precomputed Bool "
978
     "mask is true."},
979
    {"reverse", LazyFrame_reverse, METH_NOARGS,
980
     "reverse() -> LazyFrame with rows in reverse order."},
981
    {"sort_by_multi", DFTU_PYCFUNCTION(LazyFrame_sort_by_multi),
982
     METH_VARARGS | METH_KEYWORDS,
983
     "sort_by_multi(names, descending=False) -> LazyFrame stably sorted "
984
     "lexicographically by several key columns."},
985
    {"join", DFTU_PYCFUNCTION(LazyFrame_join), METH_VARARGS | METH_KEYWORDS,
986
     "join(other, on=None, how='inner', left_on=None, right_on=None, "
987
     "suffix='_right') -> LazyFrame hash-joined with another LazyFrame; "
988
     "how=inner|left|right|outer|full|semi|anti|cross."},
989
    {"concat", LazyFrame_concat, METH_VARARGS,
990
     "concat(*others) -> LazyFrame with every row of this plan followed by "
991
     "every row of each other plan in turn; schemas must match."},
992
    {"unique", DFTU_PYCFUNCTION(LazyFrame_unique), METH_VARARGS | METH_KEYWORDS,
993
     "unique(subset=None) -> LazyFrame with duplicate rows removed (keep "
994
     "first), keyed on every column or on the subset names."},
995
    {"drop_duplicates", DFTU_PYCFUNCTION(LazyFrame_unique),
996
     METH_VARARGS | METH_KEYWORDS,
997
     "drop_duplicates(subset=None) -> alias of unique."},
998
    {"sample", DFTU_PYCFUNCTION(LazyFrame_sample), METH_VARARGS | METH_KEYWORDS,
999
     "sample(n, seed=0) -> LazyFrame deterministic n-row sample."},
1000
    {"is_duplicated", LazyFrame_is_duplicated, METH_NOARGS,
1001
     "is_duplicated() -> LazyFrame Bool column, true where the row repeats."},
1002
    {"is_unique", LazyFrame_is_unique, METH_NOARGS,
1003
     "is_unique() -> LazyFrame Bool column, true where the row is unique."},
1004
    {"group_by_dynamic", DFTU_PYCFUNCTION(LazyFrame_group_by_dynamic),
1005
     METH_VARARGS | METH_KEYWORDS,
1006
     "group_by_dynamic(time_col, every, period=None, aggs=[], origin=0, "
1007
     "origin_min=False) -> LazyFrame of time-window aggregates."},
1008
    {"pivot", DFTU_PYCFUNCTION(LazyFrame_pivot), METH_VARARGS | METH_KEYWORDS,
1009
     "pivot(index, on, values, agg='first') -> LazyFrame reshaped long to "
1010
     "wide."},
1011
    {"to_dummies", LazyFrame_to_dummies, METH_O,
1012
     "to_dummies(column) -> LazyFrame one-hot encoding of the column."},
1013
    {"describe", LazyFrame_describe, METH_NOARGS,
1014
     "describe() -> LazyFrame of per-column summary statistics."},
1015
    {"memory_budget", LazyFrame_memory_budget, METH_O,
1016
     "memory_budget(bytes) -> LazyFrame with the out-of-core spill budget "
1017
     "set."},
1018
    {"auto_spill", LazyFrame_auto_spill, METH_NOARGS,
1019
     "auto_spill() -> LazyFrame with the spill budget at ~1/3 of memory."},
1020
    {"output_schema", LazyFrame_output_schema, METH_NOARGS,
1021
     "output_schema() -> [(name, dtype id)] without running; a type the "
1022
     "plan cannot know statically is 0 (unknown)."},
1023
    {"stream", DFTU_PYCFUNCTION(LazyFrame_stream), METH_VARARGS | METH_KEYWORDS,
1024
     "stream(morsel_rows=0, runtime=None) -> iterator of DataFrame chunks."},
1025
    {"schema", LazyFrame_schema, METH_NOARGS,
1026
     "schema() -> list[str] of output column names, or [] if data-dependent."},
1027
    {"explain", LazyFrame_explain, METH_NOARGS,
1028
     "explain() -> str of the optimized plan."},
1029
    {"collect", DFTU_PYCFUNCTION(LazyFrame_collect),
1030
     METH_VARARGS | METH_KEYWORDS,
1031
     "collect(morsel_rows=65536) -> _DataFrame, running the pipeline."},
1032
    {nullptr, nullptr, 0, nullptr}};
1033

1034
PyObject* collect_all_py(PyObject*, PyObject* args, PyObject* kwds) {
128 ✔
1035
    PyObject* arg = nullptr;
128 ✔
1036
    PyObject* runtime_arg = nullptr;
128 ✔
1037
    static const char* kw[] = {"plans", "runtime", nullptr};
1038
    if (!PyArg_ParseTupleAndKeywords(args, kwds, "O|O", const_cast<char**>(kw),
128 !
1039
                                     &arg, &runtime_arg))
NEW
1040
        return nullptr;
×
1041
    std::shared_ptr<dftracer::utils::Runtime> rt =
1042
        runtime_from_arg(runtime_arg);
128 !
1043
    if (!rt) return nullptr;
128 ✔
1044
    PyObject* seq = PySequence_Fast(arg, "collect_all expects a sequence");
128 !
1045
    if (!seq) return nullptr;
128 ✔
1046
    const Py_ssize_t n = PySequence_Fast_GET_SIZE(seq);
128 !
1047
    std::vector<LazyFrame> plans;
128 ✔
1048
    plans.reserve(static_cast<std::size_t>(n));
128 !
1049
    for (Py_ssize_t i = 0; i < n; ++i) {
314 ✔
1050
        LazyFrameObject* b = as_lazyframe(PySequence_Fast_GET_ITEM(seq, i));
186 !
1051
        if (!b) {
186 !
1052
            Py_DECREF(seq);
×
NEW
1053
            return nullptr;
×
1054
        }
1055
        plans.push_back(b->lf);
186 !
1056
    }
93 ✔
1057
    Py_DECREF(seq);
64 !
1058
    std::vector<dataframe::DataFrame> frames;
128 ✔
1059
    if (!run_blocking([&] {
192 !
1060
            frames = rt->submit(dataframe::collect_all(std::move(plans))).get();
128 !
1061
        }))
128 ✔
NEW
1062
        return nullptr;
×
1063
    PyObject* out = PyList_New(static_cast<Py_ssize_t>(frames.size()));
128 !
1064
    if (!out) return nullptr;
128 ✔
1065
    for (std::size_t i = 0; i < frames.size(); ++i) {
314 ✔
1066
        PyObject* df =
93 ✔
1067
            dftracer::utils::python::wrap_dataframe(std::move(frames[i]));
186 !
1068
        if (!df) {
186 !
1069
            Py_DECREF(out);
×
NEW
1070
            return nullptr;
×
1071
        }
1072
        PyList_SET_ITEM(out, static_cast<Py_ssize_t>(i), df);
186 !
1073
    }
93 ✔
1074
    return out;
128 ✔
1075
}
128 ✔
1076

1077
PyMethodDef lazyframe_module_methods[] = {
1078
    {"collect_all", DFTU_PYCFUNCTION(collect_all_py),
1079
     METH_VARARGS | METH_KEYWORDS,
1080
     "collect_all(plans) -> list of DataFrames, one per plan in order; plans "
1081
     "over the same trace base share one scan."},
1082
    {nullptr, nullptr, 0, nullptr}};
1083

1084
}  // namespace
1085

1086
namespace dftracer::utils::python {
1087

1088
int init_lazyframe(PyObject* m) {
2 ✔
1089
    LazyFrameType = {};
2 ✔
1090
    LazyFrameType.ob_base = {PyObject_HEAD_INIT(nullptr) 0};
2 ✔
1091
    LazyFrameType.tp_name = "dftracer_utils_ext._LazyFrame";
2 ✔
1092
    LazyFrameType.tp_basicsize = sizeof(LazyFrameObject);
2 ✔
1093
    LazyFrameType.tp_flags = Py_TPFLAGS_DEFAULT;
2 ✔
1094
    LazyFrameType.tp_doc = "A deferred query over a native vec batch.";
2 ✔
1095
    LazyFrameType.tp_dealloc = reinterpret_cast<destructor>(LazyFrame_dealloc);
2 ✔
1096
    LazyFrameType.tp_methods = LazyFrame_methods;
2 ✔
1097
    LazyFrameType.tp_new = nullptr;  // created only by DataFrame.lazy()
2 ✔
1098
    if (register_type(m, &LazyFrameType, "_LazyFrame") < 0) return -1;
2 ✔
1099
    return PyModule_AddFunctions(m, lazyframe_module_methods);
2 ✔
1100
}
1 ✔
1101

1102
PyObject* wrap_lazyframe(dataframe::LazyFrame&& lf) {
2,098 ✔
1103
    return make_lazyframe(std::move(lf));
2,098 ✔
1104
}
1105

1106
const dataframe::LazyFrame* lazyframe_of(PyObject* o) {
56 ✔
1107
    LazyFrameObject* b = as_lazyframe(o);
56 ✔
1108
    return b ? &b->lf : nullptr;
56 !
1109
}
1110

1111
}  // namespace dftracer::utils::python
1112

1113
#else   // !DFTRACER_UTILS_ENABLE_ARROW
1114

1115
namespace dftracer::utils::python {
1116
int init_lazyframe(PyObject*) { return 0; }
1117
PyObject* wrap_lazyframe(dftracer::utils::dataframe::LazyFrame&&) {
1118
    PyErr_SetString(PyExc_RuntimeError, "LazyFrame requires the Arrow build");
1119
    return nullptr;
1120
}
1121

1122
const dftracer::utils::dataframe::LazyFrame* lazyframe_of(PyObject*) {
1123
    PyErr_SetString(PyExc_RuntimeError, "LazyFrame requires the Arrow build");
1124
    return nullptr;
1125
}
1126
}  // namespace dftracer::utils::python
1127

1128
#endif  // DFTRACER_UTILS_ENABLE_ARROW
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