join iterate

This commit is contained in:
John Kerl 2015-07-09 18:10:32 -04:00
parent c5f4eba8b6
commit 677e636893

View file

@ -8,23 +8,7 @@
#include "input/lrec_readers.h"
#include "cli/argparse.h"
// xxx comment
#define OPTION_UNSPECIFIED ((char)0xff)
#define LEFT_STATE_0_PREFILL 0
#define LEFT_STATE_1_FULL 1
#define LEFT_STATE_2_LAST_BUCKET 2
#define LEFT_STATE_3_EOF 3
// ----------------------------------------------------------------
// xxx left state:
// (0) pre-fill: Lv == null, peek == null, leof = false
// (1) midstream: Lv != null, peek != null, leof = false
// (2) last bucket: Lv != null, peek == null, leof = true
// (3) leof: Lv == null, peek == null, leof = true
// ----------------------------------------------------------------
// ----------------------------------------------------------------
#define OPTION_UNSPECIFIED ((char)0xff) // xxx comment
typedef struct _join_bucket_t {
slls_t* pleft_field_values;
@ -32,18 +16,7 @@ typedef struct _join_bucket_t {
int was_paired;
} join_bucket_t;
typedef struct _join_bucket_keeper_t {
lrec_reader_t* plrec_reader;
void* pvhandle;
context_t* pctx;
int state;
slls_t* pleft_field_values;
sllv_t* precords;
lrec_t* prec_peek;
int leof;
} join_bucket_keeper_t;
// ----------------------------------------------------------------
typedef struct _mapper_join_opts_t {
// xxx prefix for left non-join field names
@ -77,7 +50,7 @@ typedef struct _mapper_join_state_t {
hss_t* pright_field_name_set;
// xxx cmt for sorted
join_bucket_keeper_t* pjoin_bucket_keeper;
//xxx join_bucket_keeper_t* pjoin_bucket_keeper;
// xxx key_field -> join_field (or left_field?) thruout
lhmslv_t* pbuckets_by_key_field_names; // For unsorted input
@ -89,273 +62,9 @@ typedef struct _mapper_join_state_t {
static void merge_options(mapper_join_opts_t* popts);
static void ingest_left_file(mapper_join_state_t* pstate);
// ----------------------------------------------------------------
static join_bucket_keeper_t* join_bucket_keeper_alloc(mapper_join_opts_t* popts) {
join_bucket_keeper_t* pkeeper = mlr_malloc_or_die(sizeof(join_bucket_keeper_t));
merge_options(popts);
pkeeper->plrec_reader = lrec_reader_alloc(popts->input_file_format, popts->use_mmap_for_read,
popts->irs, popts->ifs, popts->allow_repeat_ifs, popts->ips, popts->allow_repeat_ips);
pkeeper->pvhandle = pkeeper->plrec_reader->popen_func(popts->left_file_name);
pkeeper->plrec_reader->psof_func(pkeeper->plrec_reader->pvstate);
pkeeper->pctx = mlr_malloc_or_die(sizeof(context_t));
pkeeper->pctx->nr = 0; // xxx make an init func & use it here & in stream.c?
pkeeper->pctx->fnr = 0; // xxx incr this in the readers ...
pkeeper->pctx->filenum = 1;
pkeeper->pctx->filename = popts->left_file_name;
pkeeper->pleft_field_values = NULL;
pkeeper->precords = sllv_alloc();
pkeeper->prec_peek = NULL;
pkeeper->leof = FALSE;
pkeeper->state = LEFT_STATE_0_PREFILL;
return pkeeper;
}
// ----------------------------------------------------------------
static void join_bucket_keeper_free(join_bucket_keeper_t* pkeeper) {
if (pkeeper->pleft_field_values != NULL)
slls_free(pkeeper->pleft_field_values);
if (pkeeper->precords != NULL)
sllv_free(pkeeper->precords);
free(pkeeper);
pkeeper->plrec_reader->pclose_func(pkeeper->pvhandle);
}
// ----------------------------------------------------------------
// xxx left state:
// (0) pre-fill: Lv == null, peek == null, leof = false
// (1) midstream: Lv != null, peek != null, leof = false
// (2) last bucket: Lv != null, peek == null, leof = true
// (3) leof: Lv == null, peek == null, leof = true
static int join_bucket_keeper_get_state(join_bucket_keeper_t* pkeeper) {
if (pkeeper->pleft_field_values == NULL) {
if (pkeeper->leof)
return LEFT_STATE_3_EOF;
else
return LEFT_STATE_0_PREFILL;
} else {
if (pkeeper->prec_peek == NULL)
return LEFT_STATE_2_LAST_BUCKET;
else
return LEFT_STATE_1_FULL;
}
}
// xxx put bucket & bucket-keeper into separate files with separate UTs
static void join_bucket_keeper_initial_fill(join_bucket_keeper_t* pkeeper, slls_t* pleft_field_names) {
pkeeper->prec_peek = pkeeper->plrec_reader->pprocess_func(pkeeper->pvhandle,
pkeeper->plrec_reader->pvstate, pkeeper->pctx);
if (pkeeper->prec_peek == NULL) {
pkeeper->leof = TRUE;
return;
}
pkeeper->pleft_field_values = mlr_selected_values_from_record(pkeeper->prec_peek,
pleft_field_names);
sllv_add(pkeeper->precords, pkeeper->prec_peek);
pkeeper->prec_peek = NULL;
while (TRUE) {
pkeeper->prec_peek = pkeeper->plrec_reader->pprocess_func(pkeeper->pvhandle,
pkeeper->plrec_reader->pvstate, pkeeper->pctx);
if (pkeeper->prec_peek == NULL) {
pkeeper->leof = TRUE;
break;
}
// xxx make a function to compare w/o copy
slls_t* pnext_field_values = mlr_selected_values_from_record(pkeeper->prec_peek,
pleft_field_names);
int cmp = slls_compare_lexically(pkeeper->pleft_field_values, pnext_field_values);
if (cmp != 0) {
break;
}
sllv_add(pkeeper->precords, pkeeper->prec_peek);
pkeeper->prec_peek = NULL;
}
}
// if (pkeeper->leof) {
// int cmp = slls_compare_lexically(pkeeper->pbucket->pjoin_values, pright_field_values);
// // xxx think through various cases
// if (/* xxx stub */ cmp == 999) {
// // rename: "bucket" used at different nesting levels. it's confusing.
// *ppbucket_paired = pkeeper->precords;
// } else {
// *ppbucket_left_unpaired = pkeeper->precords;
// }
// pkeeper->pbucket = NULL;
// return;
// }
//
// // xxx now that we've got a peek record: fill the bucket with like records
// // until there's a non-like on-deck.
// //
// // xxx do this only if it's time for a change.
// //
// // xxx rename "peek" to "on-deck"?
//
//#if 0
// sllv_empty(pkeeper->pbucket);
// while (TRUE) {
// sllv_add(pkeeper->pbucket, pkeeper->prec_peek);
// pkeeper->prec_peek = pkeeper->plrec_reader->pprocess_func(pkeeper->pvhandle,
// pkeeper->plrec_reader->pvstate, pkeeper->pctx);
// get selected keys
// if (pkeeper->prec_peek == NULL)
// pkeeper->leof = TRUE;
// x
// }
//#endif
//
// // xxx stub
// lrec_t* pleft_rec = pkeeper->plrec_reader->pprocess_func(pkeeper->pvhandle, pkeeper->plrec_reader->pvstate,
// pkeeper->pctx);
// sllv_t* pfoo = sllv_alloc();
// if (pleft_rec != NULL)
// sllv_add(pfoo, pleft_rec);
// *ppbucket_paired = pfoo;
// xxx cmt re who frees
static void join_bucket_keeper_emit(join_bucket_keeper_t* pkeeper,
slls_t* pleft_field_names, slls_t* pright_field_values,
sllv_t** ppbucket_paired, sllv_t** ppbucket_left_unpaired)
{
*ppbucket_paired = NULL;
*ppbucket_left_unpaired = NULL;
int cmp = 0;
if (pkeeper->state == LEFT_STATE_0_PREFILL) {
// try fill Lv & peek; next state is 1,2,3 & continue from there.
join_bucket_keeper_initial_fill(pkeeper, pleft_field_names);
pkeeper->state = join_bucket_keeper_get_state(pkeeper);
}
//
//typedef struct _join_bucket_keeper_t {
// lrec_reader_t* plrec_reader;
// void* pvhandle;
// context_t* pctx;
//
// int state;
// slls_t* pleft_field_values;
// sllv_t* precords;
// lrec_t* prec_peek;
// int leof;
//} join_bucket_keeper_t;
switch (pkeeper->state) {
case LEFT_STATE_1_FULL:
case LEFT_STATE_2_LAST_BUCKET: // Intentional fall-through
cmp = slls_compare_lexically(pkeeper->pleft_field_values, pright_field_values);
if (cmp < 0) {
*ppbucket_left_unpaired = pkeeper->precords;
// xxx advance left ...
} else if (cmp == 0) {
*ppbucket_paired = pkeeper->precords;
}
break;
case LEFT_STATE_3_EOF:
break;
default:
fprintf(stderr, "%s: internal coding error: failed transition from prefill state.\n",
MLR_GLOBALS.argv0);
exit(1);
break;
}
pkeeper->state = join_bucket_keeper_get_state(pkeeper);
// // xxx cmt why rec-peek
// if (!pkeeper->leof && pkeeper->prec_peek == NULL) {
// pkeeper->prec_peek = pkeeper->plrec_reader->pprocess_func(pkeeper->pvhandle,
// pkeeper->plrec_reader->pvstate, pkeeper->pctx);
// if (pkeeper->prec_peek == NULL)
// pkeeper->leof = TRUE;
// }
//
// // xxx not quite: only if *already* leof.
// if (pkeeper->leof) {
// int cmp = slls_compare_lexically(pkeeper->pbucket->pjoin_values, pright_field_values);
// // xxx think through various cases
// if (/* xxx stub */ cmp == 999) {
// // rename: "bucket" used at different nesting levels. it's confusing.
// *ppbucket_paired = pkeeper->precords;
// } else {
// *ppbucket_left_unpaired = pkeeper->precords;
// }
// pkeeper->pbucket = NULL;
// return;
// }
//
// // xxx now that we've got a peek record: fill the bucket with like records
// // until there's a non-like on-deck.
// //
// // xxx do this only if it's time for a change.
// //
// // xxx rename "peek" to "on-deck"?
//
//#if 0
// sllv_empty(pkeeper->pbucket);
// while (TRUE) {
// sllv_add(pkeeper->pbucket, pkeeper->prec_peek);
// pkeeper->prec_peek = pkeeper->plrec_reader->pprocess_func(pkeeper->pvhandle,
// pkeeper->plrec_reader->pvstate, pkeeper->pctx);
// get selected keys
// if (pkeeper->prec_peek == NULL)
// pkeeper->leof = TRUE;
// x
// }
//#endif
//
// // xxx stub
// lrec_t* pleft_rec = pkeeper->plrec_reader->pprocess_func(pkeeper->pvhandle, pkeeper->plrec_reader->pvstate,
// pkeeper->pctx);
// sllv_t* pfoo = sllv_alloc();
// if (pleft_rec != NULL)
// sllv_add(pfoo, pleft_rec);
// *ppbucket_paired = pfoo;
}
// xxx need a drain-hook for returning the final left-unpaireds after right EOF.
// ----------------------------------------------------------------
// +-----------+-----------+-----------+-----------+-----------+-----------+
// | L R | L R | L R | L R | L R | L R |
// + --- --- + --- --- + --- --- + --- --- + --- --- + --- --- +
// | a | a | e | a | e e | e e |
// | b | e | e | e e | e | e e |
// | e | e | e | e | e | e |
// | e | e | f | e | f | g g |
// | e | f | g | g | g | g |
// | g | g | g | g | g | |
// | g | g | h | | | |
// +-----------+-----------+-----------+-----------+-----------+-----------+
// Cases:
// * 1st emit, right row < 1st left row
// * 1st emit, right row == 1st left row
// * 1st emit, right row > 1st left row
// * subsequent emit, right row < 1st left row
// * subsequent emit, right row == 1st left row
// * subsequent emit, right row > 1st left row
// * new left EOF, right row < 1st left row
// * new left EOF, right row == 1st left row
// * new left EOF, right row > 1st left row
// * old left EOF, right row < 1st left row
// * old left EOF, right row == 1st left row
// * old left EOF, right row > 1st left row
// ----------------------------------------------------------------
static sllv_t* mapper_join_process_sorted(lrec_t* pright_rec, context_t* pctx, void* pvstate) {
#if 0
mapper_join_state_t* pstate = (mapper_join_state_t*)pvstate;
// This can't be done in the CLI-parser since it requires information which
@ -430,6 +139,9 @@ static sllv_t* mapper_join_process_sorted(lrec_t* pright_rec, context_t* pctx, v
}
return pout_recs;
#else
return NULL;
#endif
}
// ----------------------------------------------------------------
@ -514,8 +226,9 @@ static void mapper_join_free(void* pvstate) {
if (pstate->popts->poutput_field_names != NULL)
slls_free(pstate->popts->poutput_field_names);
if (pstate->pjoin_bucket_keeper != NULL)
join_bucket_keeper_free(pstate->pjoin_bucket_keeper);
// xxx
// if (pstate->pjoin_bucket_keeper != NULL)
// join_bucket_keeper_free(pstate->pjoin_bucket_keeper);
}
// ----------------------------------------------------------------
@ -590,7 +303,7 @@ static mapper_t* mapper_join_alloc(mapper_join_opts_t* popts)
pstate->pleft_field_name_set = hss_from_slls(popts->pleft_field_names);
pstate->pright_field_name_set = hss_from_slls(popts->pright_field_names);
pstate->pbuckets_by_key_field_names = NULL;
pstate->pjoin_bucket_keeper = NULL;
// xxx pstate->pjoin_bucket_keeper = NULL;
pmapper->pvstate = (void*)pstate;
if (popts->allow_unsorted_input) {