quantiles checkpoint

This commit is contained in:
John Kerl 2015-05-10 21:35:02 -04:00
parent 4a96f9a283
commit c3d1ceecd4

View file

@ -16,7 +16,7 @@
// ================================================================
typedef void acc_dingest_func_t(void* pvstate, double val);
typedef void acc_singest_func_t(void* pvstate, char* val);
typedef void acc_emit_func_t(void* pvstate, char* value_field_name, lrec_t* poutrec);
typedef void acc_emit_func_t(void* pvstate, char* value_field_name, char* acc_name, lrec_t* poutrec);
typedef struct _acc_t {
void* pvstate;
@ -36,9 +36,9 @@ void acc_count_singest(void* pvstate, char* val) {
acc_count_state_t* pstate = pvstate;
pstate->count++;
}
void acc_count_emit(void* pvstate, char* value_field_name, lrec_t* poutrec) {
void acc_count_emit(void* pvstate, char* value_field_name, char* acc_name, lrec_t* poutrec) {
acc_count_state_t* pstate = pvstate;
char* key = mlr_paste_2_strings(value_field_name, "_count");
char* key = mlr_paste_3_strings(value_field_name, "_", acc_name);
char* val = mlr_alloc_string_from_ull(pstate->count);
lrec_put(poutrec, key, val, LREC_FREE_ENTRY_KEY|LREC_FREE_ENTRY_VALUE);
}
@ -63,9 +63,9 @@ void acc_sum_dingest(void* pvstate, double val) {
acc_sum_state_t* pstate = pvstate;
pstate->sum += val;
}
void acc_sum_emit(void* pvstate, char* value_field_name, lrec_t* poutrec) {
void acc_sum_emit(void* pvstate, char* value_field_name, char* acc_name, lrec_t* poutrec) {
acc_sum_state_t* pstate = pvstate;
char* key = mlr_paste_2_strings(value_field_name, "_sum");
char* key = mlr_paste_3_strings(value_field_name, "_", acc_name);
char* val = mlr_alloc_string_from_double(pstate->sum, pstate->pstatx->ofmt);
lrec_put(poutrec, key, val, LREC_FREE_ENTRY_KEY|LREC_FREE_ENTRY_VALUE);
}
@ -92,10 +92,10 @@ void acc_avg_dingest(void* pvstate, double val) {
pstate->sum += val;
pstate->count++;
}
void acc_avg_emit(void* pvstate, char* value_field_name, lrec_t* poutrec) {
void acc_avg_emit(void* pvstate, char* value_field_name, char* acc_name, lrec_t* poutrec) {
acc_avg_state_t* pstate = pvstate;
double quot = pstate->sum / pstate->count;
char* key = mlr_paste_2_strings(value_field_name, "_avg");
char* key = mlr_paste_3_strings(value_field_name, "_", acc_name);
char* val = mlr_alloc_string_from_double(quot, pstate->pstatx->ofmt);
lrec_put(poutrec, key, val, LREC_FREE_ENTRY_KEY|LREC_FREE_ENTRY_VALUE);
}
@ -126,9 +126,9 @@ void acc_stddev_avgeb_dingest(void* pvstate, double val) {
pstate->sumx += val;
pstate->sumx2 += val*val;
}
void acc_stddev_avgeb_emit(void* pvstate, char* value_field_name, lrec_t* poutrec) {
void acc_stddev_avgeb_emit(void* pvstate, char* value_field_name, char* acc_name, lrec_t* poutrec) {
acc_stddev_avgeb_state_t* pstate = pvstate;
char* key = mlr_paste_2_strings(value_field_name, pstate->do_avgeb ? "_avgeb" : "_stddev");
char* key = mlr_paste_3_strings(value_field_name, "_", acc_name);
if (pstate->count < 2LL) {
lrec_put(poutrec, key, "", LREC_FREE_ENTRY_KEY);
} else {
@ -177,9 +177,9 @@ void acc_min_dingest(void* pvstate, double val) {
pstate->min = val;
}
}
void acc_min_emit(void* pvstate, char* value_field_name, lrec_t* poutrec) {
void acc_min_emit(void* pvstate, char* value_field_name, char* acc_name, lrec_t* poutrec) {
acc_min_state_t* pstate = pvstate;
char* key = mlr_paste_2_strings(value_field_name, "_min");
char* key = mlr_paste_3_strings(value_field_name, "_", acc_name);
if (pstate->have_min) {
char* val = mlr_alloc_string_from_double(pstate->min, pstate->pstatx->ofmt);
lrec_put(poutrec, key, val, LREC_FREE_ENTRY_KEY|LREC_FREE_ENTRY_VALUE);
@ -216,9 +216,9 @@ void acc_max_dingest(void* pvstate, double val) {
pstate->max = val;
}
}
void acc_max_emit(void* pvstate, char* value_field_name, lrec_t* poutrec) {
void acc_max_emit(void* pvstate, char* value_field_name, char* acc_name, lrec_t* poutrec) {
acc_max_state_t* pstate = pvstate;
char* key = mlr_paste_2_strings(value_field_name, "_max");
char* key = mlr_paste_3_strings(value_field_name, "_", acc_name);
if (pstate->have_max) {
char* val = mlr_alloc_string_from_double(pstate->max, pstate->pstatx->ofmt);
lrec_put(poutrec, key, val, LREC_FREE_ENTRY_KEY|LREC_FREE_ENTRY_VALUE);
@ -255,7 +255,7 @@ void acc_mode_singest(void* pvstate, char* val) {
pe->value++;
}
}
void acc_mode_emit(void* pvstate, char* value_field_name, lrec_t* poutrec) {
void acc_mode_emit(void* pvstate, char* value_field_name, char* acc_name, lrec_t* poutrec) {
acc_mode_state_t* pstate = pvstate;
int max_count = 0;
char* max_key = "";
@ -266,7 +266,7 @@ void acc_mode_emit(void* pvstate, char* value_field_name, lrec_t* poutrec) {
max_count = count;
}
}
char* key = mlr_paste_2_strings(value_field_name, "_mode");
char* key = mlr_paste_3_strings(value_field_name, "_", acc_name);
lrec_put(poutrec, key, max_key, LREC_FREE_ENTRY_KEY);
}
// xxx somewhere note a crucial assumption: while pctx is passed in
@ -289,49 +289,6 @@ acc_t* acc_mode_alloc(static_context_t* pstatx) {
return pacc;
}
// ----------------------------------------------------------------
typedef struct _acc_foo_state_t {
int have_foo;
double foo;
static_context_t* pstatx;
} acc_foo_state_t;
void acc_foo_dingest(void* pvstate, double val) {
acc_foo_state_t* pstate = pvstate;
if (pstate->have_foo) {
if (val > pstate->foo)
pstate->foo = val;
} else {
pstate->have_foo = TRUE;
pstate->foo = val;
}
}
void acc_foo_emit(void* pvstate, char* value_field_name, lrec_t* poutrec) {
acc_foo_state_t* pstate = pvstate;
char* key1 = mlr_paste_2_strings(value_field_name, "_foo");
char* key2 = mlr_paste_2_strings(value_field_name, "_bar");
if (pstate->have_foo) {
char* val = mlr_alloc_string_from_double(pstate->foo, pstate->pstatx->ofmt);
lrec_put(poutrec, key1, val, LREC_FREE_ENTRY_KEY|LREC_FREE_ENTRY_VALUE);
val = mlr_alloc_string_from_double(-pstate->foo, pstate->pstatx->ofmt);
lrec_put(poutrec, key2, val, LREC_FREE_ENTRY_KEY|LREC_FREE_ENTRY_VALUE);
} else {
lrec_put(poutrec, key1, "", LREC_FREE_ENTRY_KEY);
lrec_put(poutrec, key2, "", LREC_FREE_ENTRY_KEY);
}
}
acc_t* acc_foo_alloc(static_context_t* pstatx) {
acc_t* pacc = mlr_malloc_or_die(sizeof(acc_t));
acc_foo_state_t* pstate = mlr_malloc_or_die(sizeof(acc_foo_state_t));
pstate->have_foo = FALSE;
pstate->foo = -999.0;
pstate->pstatx = pstatx;
pacc->pvstate = (void*)pstate;
pacc->psingest_func = NULL;
pacc->pdingest_func = &acc_foo_dingest;
pacc->pemit_func = &acc_foo_emit;
return pacc;
}
// ----------------------------------------------------------------
typedef struct _acc_lookup_t {
char* name;
@ -346,7 +303,6 @@ static acc_lookup_t acc_lookup_table[] = {
{"min", acc_min_alloc},
{"max", acc_max_alloc},
{"mode", acc_mode_alloc},
{"foo", acc_foo_alloc},
};
static int acc_lookup_table_length = sizeof(acc_lookup_table) / sizeof(acc_lookup_table[0]);
@ -534,9 +490,8 @@ static sllv_t* mapper_stats1_emit(mapper_stats1_state_t* pstate) {
sllv_t* poutrecs = sllv_alloc();
for (lhmslve_t* pa = pstate->groups->phead; pa != NULL; pa = pa->pnext) {
lrec_t* poutrec = lrec_unbacked_alloc();
slls_t* pgroup_by_field_values = pa->key;
lrec_t* poutrec = lrec_unbacked_alloc();
// Add in a=s,b=t fields:
sllse_t* pb = pstate->pgroup_by_field_names->phead;
@ -551,15 +506,21 @@ static sllv_t* mapper_stats1_emit(mapper_stats1_state_t* pstate) {
for (lhmsve_t* pd = group_to_acc_field->phead; pd != NULL; pd = pd->pnext) {
char* value_field_name = pd->key;
lhmsv_t* acc_field_to_acc_state = pd->value;
// for "count", "sum"
for (lhmsve_t* pe = acc_field_to_acc_state->phead; pe != NULL; pe = pe->pnext) {
if (streq(pe->key, fake_acc_name_for_setups))
for (sllse_t* pe = pstate->paccumulator_names->phead; pe != NULL; pe = pe->pnext) {
char* acc_name = pe->value;
if (streq(acc_name, fake_acc_name_for_setups))
continue;
acc_t* pacc = pe->value;
pacc->pemit_func(pacc->pvstate, value_field_name, poutrec);
acc_t* pacc = lhmsv_get(acc_field_to_acc_state, acc_name);
if (pacc == NULL) {
// xxx needs argv0 from mlr_globals
fprintf(stderr, "mlr stats1: internal coding error: acc_name \"%s\" has gone missing.\n",
acc_name);
exit(1);
}
pacc->pemit_func(pacc->pvstate, value_field_name, acc_name, poutrec);
}
}
sllv_add(poutrecs, poutrec);
}
sllv_add(poutrecs, NULL);