Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 11 additions & 3 deletions src/mem/heap.c
Original file line number Diff line number Diff line change
Expand Up @@ -1683,8 +1683,11 @@ void ray_free(ray_t* v) {
* pool before appending — has to ask the registry, and only a
* string column was ever registered. So the lock stays off the
* ordinary free entirely, and a mutated column pays it once. */
/* A column loaded with an inline index registered its region too
* (col.c): the index may have been detached since, and only the
* descriptor still knows the mapped length. */
ray_file_map_t* m = col_map;
if (!m && v->type == RAY_STR) m = ray_file_map_lookup(v);
if (!m) m = ray_file_map_lookup(v);
if (m) {
ray_file_map_release(m);
if (h) RAY_STAT(h->stats.free_count++);
Expand Down Expand Up @@ -2942,9 +2945,14 @@ void ray_parallel_end(void) {
* else is the page-rounded payload plus an inline passenger index. */
static size_t mapped_block_bytes(const ray_t* v) {
if (v->type == RAY_TABLE || v->type == RAY_DICT || v->type == RAY_LIST) return 0;
if (v->type == RAY_STR) {
{
/* Same route as ray_free: the pool's descriptor for a string
* column, else the registry — which also holds every indexed
* mapped column (col.c), whether or not its index is still
* attached. */
ray_file_map_t* m = NULL;
if (v->str_pool && !RAY_IS_ERR(v->str_pool) && v->str_pool->mmod == 3)
if (v->type == RAY_STR && v->str_pool && !RAY_IS_ERR(v->str_pool) &&
v->str_pool->mmod == 3)
m = v->str_pool->file_map;
if (!m) m = ray_file_map_lookup(v);
if (m) return m->len;
Expand Down
91 changes: 60 additions & 31 deletions src/ops/agg.c
Original file line number Diff line number Diff line change
Expand Up @@ -169,12 +169,14 @@ static ray_t* agg_parted_avg(ray_t* x) {
if (!agg_parted_numeric_base(base)) return ray_error("type", "avg expects a numeric or temporal parted column, got %s", ray_type_name(base));
ray_t** segs = (ray_t**)ray_data(x);
double sum = 0.0;
int64_t hi = 0; uint64_t lo = 0; /* integer segments: exact 128-bit sum */
int64_t cnt = 0;
bool fp = base == RAY_F64 || base == RAY_F32;
for (int64_t s = 0; s < x->len; s++) {
ray_t* seg = segs[s];
if (!seg) continue;
int has_nulls = ray_vec_may_have_nulls(seg);
if (base == RAY_F64 || base == RAY_F32) {
if (fp) {
for (int64_t i = 0; i < seg->len; i++) {
if (has_nulls && ray_vec_is_null(seg, i)) continue;
if (base == RAY_F64) sum += ((double*)ray_data(seg))[i];
Expand All @@ -184,12 +186,12 @@ static ray_t* agg_parted_avg(ray_t* x) {
} else {
for (int64_t i = 0; i < seg->len; i++) {
if (has_nulls && ray_vec_is_null(seg, i)) continue;
sum += (double)agg_read_i64(seg, i); cnt++;
ray_i128_add(&hi, &lo, agg_read_i64(seg, i)); cnt++;
}
}
}
if (cnt == 0) return ray_typed_null(-RAY_F64);
return make_f64(sum / (double)cnt);
return make_f64((fp ? sum : ray_i128_to_f64(hi, lo)) / (double)cnt);
}

static ray_t* agg_parted_prod(ray_t* x) {
Expand Down Expand Up @@ -404,36 +406,63 @@ static ray_t* agg_pair_vec(ray_t* x, ray_t* y, uint16_t op) {
return ray_f64(ray_f64_fin(num / sqrt(dx * dy)));
}

/* Whole-column sum / non-null count of an integer column from its
* chunk-zone index (per-chunk int64 sums with wraparound and non-null
* counts), in O(n_chunks). `exact_f64` reports whether every partial sum
* an accumulation in double could meet stays below 2^53 in magnitude, i.e.
* the sum converted to double equals any double accumulation of the rows.
* Returns false when the column has no such index for its current length. */
bool ray_zone_int_sum(ray_t* x, int64_t* sum_out, int64_t* nn_out, bool* exact_f64) {
if (!x || !ray_is_vec(x) || ray_index_kind(x) != RAY_IDX_CHUNK_ZONE) return false;
/* Per-chunk aggregates of an integer column's chunk-zone index, or NULL:
* [sum low words | non-null counts | sum high words] (the last block only
* in the current layout). */
static const int64_t* zone_aggs(ray_t* x, uint32_t* n_out, bool* have_hi) {
if (!x || !ray_is_vec(x) || ray_index_kind(x) != RAY_IDX_CHUNK_ZONE) return NULL;
ray_index_t* ix = ray_index_payload(x->index);
if (ix->built_for_len != x->len || ix->u.chunk_zone.is_f64 || !ix->u.chunk_zone.aggs)
return false;
return NULL;
uint32_t n = ix->u.chunk_zone.n_chunks;
if (ix->u.chunk_zone.aggs->len != 2 * (int64_t)n) return false;
const int64_t* ag = (const int64_t*)ray_data(ix->u.chunk_zone.aggs);
const int64_t* mins = (const int64_t*)ray_data(ix->u.chunk_zone.mins);
const int64_t* maxs = (const int64_t*)ray_data(ix->u.chunk_zone.maxs);
int64_t len = ix->u.chunk_zone.aggs->len;
if (len != 2 * (int64_t)n && len != 3 * (int64_t)n) return NULL;
*n_out = n;
*have_hi = len == 3 * (int64_t)n;
return (const int64_t*)ray_data(ix->u.chunk_zone.aggs);
}

bool ray_zone_int_sum(ray_t* x, int64_t* sum_out, int64_t* nn_out) {
uint32_t n; bool have_hi;
const int64_t* ag = zone_aggs(x, &n, &have_hi);
if (!ag) return false;
uint64_t sum = 0;
int64_t nn = 0;
double bound = 0.0;
for (uint32_t g = 0; g < n; g++) {
sum += (uint64_t)ag[g];
nn += ag[n + g];
if (ag[n + g] > 0) {
double a = fabs((double)mins[g]), b = fabs((double)maxs[g]);
bound += (a > b ? a : b) * (double)ag[n + g];
}
}
*sum_out = (int64_t)sum;
*nn_out = nn;
if (exact_f64) *exact_f64 = bound < 9007199254740992.0; /* 2^53 */
return true;
}

bool ray_zone_int_sum128(ray_t* x, int64_t* hi_out, uint64_t* lo_out, int64_t* nn_out) {
uint32_t n; bool have_hi;
const int64_t* ag = zone_aggs(x, &n, &have_hi);
if (!ag) return false;
ray_index_t* ix = ray_index_payload(x->index);
const int64_t* mins = (const int64_t*)ray_data(ix->u.chunk_zone.mins);
const int64_t* maxs = (const int64_t*)ray_data(ix->u.chunk_zone.maxs);
int64_t hi = 0; uint64_t lo = 0; int64_t nn = 0;
double bound = 0.0;
for (uint32_t g = 0; g < n; g++) {
int64_t cnt = ag[n + g];
nn += cnt;
if (have_hi) {
ray_i128_add128(&hi, &lo, ag[2 * (int64_t)n + g], (uint64_t)ag[g]);
} else {
/* no high words stored: the wrapped low word is the chunk's
* exact sum only when no partial sum could leave int64 */
if (cnt > 0) {
double a = fabs((double)mins[g]), b = fabs((double)maxs[g]);
bound += (a > b ? a : b) * (double)cnt;
}
ray_i128_add(&hi, &lo, ag[g]);
}
}
if (!have_hi && !(bound < 9223372036854775808.0)) return false; /* 2^63 */
*hi_out = hi; *lo_out = lo; *nn_out = nn;
return true;
}

Expand All @@ -453,7 +482,7 @@ ray_t* ray_sum_fn(ray_t* x) {
/* Integer columns with per-chunk sums in their zone index. */
if (x->type == RAY_I64 || x->type == RAY_I32 || x->type == RAY_I16 || x->type == RAY_U8) {
int64_t zs, zn;
if (ray_zone_int_sum(x, &zs, &zn, NULL)) return make_i64(zs);
if (ray_zone_int_sum(x, &zs, &zn)) return make_i64(zs);
}
/* Narrow/temporal types need specific return constructors that the
* DAG executor doesn't provide — use scalar path for these. TIMESTAMP
Expand Down Expand Up @@ -608,14 +637,14 @@ ray_t* ray_avg_fn(ray_t* x) {
/* Canonical admission: numeric + temporal (→ F64); SYM/STR/GUID are
* non-numeric → type error (the DAG path otherwise averaged raw ids). */
if (!agg_type_admitted(OP_AVG, x->type)) return ray_error("type", "avg expects a numeric or temporal vector, got %s", ray_type_name(x->type));
/* Integer columns with per-chunk sums: exact when no partial sum can
* leave double's integer range (then every accumulation order in
* double gives the same value). */
/* Integer columns with per-chunk sums: the exact 128-bit total is
* what the row-wise reduction computes too, so the answers agree
* bit for bit. */
if (x->type == RAY_I64 || x->type == RAY_I32 || x->type == RAY_I16 || x->type == RAY_U8) {
int64_t zs, zn; bool exact = false;
if (ray_zone_int_sum(x, &zs, &zn, &exact) && exact) {
int64_t zh, zn; uint64_t zl;
if (ray_zone_int_sum128(x, &zh, &zl, &zn)) {
if (zn == 0) return ray_typed_null(-RAY_F64);
return make_f64((double)zs / (double)zn);
return make_f64(ray_i128_to_f64(zh, zl) / (double)zn);
}
}
AGG_VEC_VIA_DAG(x, ray_avg);
Expand Down Expand Up @@ -663,7 +692,7 @@ ray_t* ray_min_fn(ray_t* x) {
* zone has them; the sentinel alone cannot tell a
* column of INT64_MAX values from an empty one. */
int64_t zs_, zn_;
bool have_nn = ray_zone_int_sum(x, &zs_, &zn_, NULL);
bool have_nn = ray_zone_int_sum(x, &zs_, &zn_);
if (have_nn ? zn_ == 0 : mn == INT64_MAX) return ray_typed_null(-x->type);
/* Preserve the column's storage width on the result. */
switch (x->type) {
Expand Down Expand Up @@ -723,7 +752,7 @@ ray_t* ray_max_fn(ray_t* x) {
* zone has them; the sentinel alone cannot tell a
* column of INT64_MIN values from an empty one. */
int64_t zs_, zn_;
bool have_nn = ray_zone_int_sum(x, &zs_, &zn_, NULL);
bool have_nn = ray_zone_int_sum(x, &zs_, &zn_);
if (have_nn ? zn_ == 0 : mx == INT64_MIN) return ray_typed_null(-x->type);
switch (x->type) {
case RAY_BOOL: return ray_bool((bool)mx);
Expand Down
8 changes: 8 additions & 0 deletions src/ops/agg_engine.c
Original file line number Diff line number Diff line change
Expand Up @@ -184,6 +184,14 @@ agg_v2_reason_t agg_v2_admission(ray_graph_t* g, ray_op_t* op, ray_t* tbl) {
const agg_vtable_t* vt = agg_resolve(ext->agg_ops[a], ic->type);
if (!vt) return AGG_V2_AGG_TYPE;
if (vt->kind != ACC_STREAMING) return AGG_V2_BUFFERED;
/* The streaming 64-bit integer avg packs its 128-bit high word with
* the group count in one word (agg_stream.c): exact while every
* group has fewer than 2^32 rows; the narrower kernels keep a plain
* int64 sum, exact below 2^31 rows. A table that large keeps the
* legacy engines, whose accumulators carry a full high word. */
if (ext->agg_ops[a] == OP_AVG && ic->type != RAY_F64 && ic->type != RAY_F32 &&
ray_table_nrows(tbl) >= ((int64_t)1 << ((ic->type == RAY_I64 || ic->type == RAY_TIMESTAMP) ? 32 : 31)))
return AGG_V2_AGG_TYPE;
}
return AGG_V2_ADMITTED;
}
Expand Down
Loading
Loading