backup for memory issues

This commit is contained in:
chhylp123
2022-05-05 19:35:06 -04:00
parent 8a2bb0b112
commit dd8888aa90
4 changed files with 173 additions and 81 deletions
+132 -62
View File
@@ -56,7 +56,7 @@ void ha_get_ul_candidates_interface(ha_abufl_t *ab, int64_t rid, char* rs, uint6
#define SEC_LEN_DIF 0.03
#define REA_ALIGN_CUTOFF 32
#define CHUNK_SIZE 500000000
#define CHUNK_SIZE 250000000
#define generic_key(x) (x)
KRADIX_SORT_INIT(gfa64, uint64_t, generic_key, 8)
@@ -399,6 +399,7 @@ typedef struct { // data structure for each step in kt_pipeline()
gdpchain_t *gdp;
// glchain_t *sec_ll;
uint64_t num_bases, num_corrected_bases, num_recorrected_bases;
int64_t n_thread;
} utepdat_t;
@@ -4941,10 +4942,12 @@ void update_ul_vec_t_ug(const ul_idx_t *uref, ul_vec_t *rch, vec_mg_lchain_t *uc
if(rch->bb.a[k].pidx == (uint32_t)-1) k = -1;
else k = rch->bb.a[k].pidx;
}
// fprintf(stderr, "-ulid->%ld\n", ulid);
rch->dd = 0;
if(sp != (uint32_t)-1) l += ep - sp;
l = (int64_t)rch->rlen - l;
// if(ulid == 1756) fprintf(stderr, "-ulid:%ld, l:%ld, rch->rlen:%u\n", ulid, l, rch->rlen);
if(l == 0) {
rch->dd = 1;
} else if(l < ((int64_t)rch->rlen)*0.001) {
@@ -5053,7 +5056,7 @@ int64_t debug_i, int64_t tid, void *km)
f = check_trans_rate_gap(&(gdp->swap), &(ll->tk), olist, hap, uref, diff_ec_ul, winLen, G_CHAIN_TRANS_RATE);
}
}
// fprintf(stderr, "(beg2) [M::%s] debug_i:%ld, qlen:%ld\n", __func__, debug_i, qlen);
// if(debug_i == 1756) fprintf(stderr, "[M::%s] ulid:%ld, qlen:%ld, f:%ld\n", __func__, debug_i, qlen, f);
if(f) update_ul_vec_t_ug(uref, rch, &(gdp->swap), debug_i);
// debug_ul_vec_t_chain(km, uref->ug->g, rch, &(gdp->dst_done), &(gdp->out));
// fprintf(stderr, "(beg3) [M::%s::tid:%ld] debug_i:%ld, qlen:%ld\n", __func__, tid, debug_i, qlen);
@@ -5080,6 +5083,35 @@ uint64_t kv_ul_ov_t_statistics(kv_ul_ov_t *olist, uint64_t qn, int64_t *occ)
return l;
}
#define kv_mem(v) ((v).m * sizeof(*((v).a)))
int64_t get_utepdat_t_mem_tid(const utepdat_t *b, int64_t tid, int64_t *mem, int64_t *mem_hab)
{
km_stat_t kmst; memset(mem, 0, sizeof(*(mem))*6);
if(b->hab) {
mem[0] += ha_ovec_mem(b->hab[tid], mem_hab);
}
if(b->ll) {
mem[1] += kv_mem(b->ll[tid].lo) + kv_mem(b->ll[tid].tk) + kv_mem(b->ll[tid].srt.a);
}
if(b->buf) {
km_stat(b->buf[tid]->km, &kmst);
mem[2] += kmst.capacity;
}
if(b->gdp) {
mem[3] += kv_mem(b->gdp[tid].l) + kv_mem(b->gdp[tid].swap) + kv_mem(b->gdp[tid].dst) + kv_mem(b->gdp[tid].out)
+ kv_mem(b->gdp[tid].path) + kv_mem(b->gdp[tid].v) + kv_mem(b->gdp[tid].f) + kv_mem(b->gdp[tid].dst_done);
}
if(b->mzs) {
mem[4] += kv_mem(b->mzs[tid]);
}
if(b->sps) {
mem[5] += kv_mem(b->sps[tid]);
}
return mem[0] + mem[1] + mem[2] + mem[3] + mem[4] + mem[5];
}
static void worker_for_ul_scall_alignment(void *data, long i, int tid) // callback for kt_for()
{
utepdat_t *s = (utepdat_t*)data;
@@ -5094,7 +5126,7 @@ static void worker_for_ul_scall_alignment(void *data, long i, int tid) // callba
// if (memcmp(UL_INF.nid.a[s->id+i].a, "d0aab024-b3a7-40fb-83cc-22c3d6d951f8", UL_INF.nid.a[s->id+i].n-1)) return;
// fprintf(stderr, "[M::%s::] ==> len: %lu\n", __func__, s->len[i]);
ha_get_ul_candidates_interface(b->abl, i, s->seq[i], s->len[i], s->opt->w, s->opt->k, s->uu, &b->olist, &b->olist_hp, &b->clist, s->opt->bw_thres,
s->opt->max_n_chain, 1, &(b->k_flag), &b->r_buf, &(b->tmp_region), NULL, &(b->sp), asm_opt.hom_cov, km);
s->opt->max_n_chain, 1, NULL/**&(b->k_flag)**/, &b->r_buf, &(b->tmp_region), NULL, &(b->sp), asm_opt.hom_cov, km);
clear_Cigar_record(&b->cigar1);
clear_Round2_alignment(&b->round2);
@@ -5146,74 +5178,85 @@ static void worker_for_ul_scall_alignment(void *data, long i, int tid) // callba
}
}
static void worker_for_ul_rescall_alignment(void *data, long i, int tid) // callback for kt_for()
{
utepdat_t *s = (utepdat_t*)data;
ha_ovec_buf_t *b = s->hab[tid];
glchain_t *bl = &(s->ll[tid]);
int64_t /**rid = s->id+i,**/ winLen = MIN((((double)THRESHOLD_MAX_SIZE)/s->opt->diff_ec_ul), WINDOW);
// uint64_t align = 0;
int fully_cov, abnormal;
// void *km = s->buf?(s->buf[tid]?s->buf[tid]->km:NULL):NULL;
// if(s->id+i!=873) return;
// fprintf(stderr, "\n[M::%s] rid:%ld, len:%lu\n", __func__, s->id+i, s->len[i]);
// if (memcmp(UL_INF.nid.a[s->id+i].a, "d0aab024-b3a7-40fb-83cc-22c3d6d951f8", UL_INF.nid.a[s->id+i].n-1)) return;
// fprintf(stderr, "[M::%s::] ==> len: %lu\n", __func__, s->len[i]);
ha_get_ul_candidates_interface(b->abl, i, s->seq[i], s->len[i], s->opt->w, s->opt->k, s->uu, &b->olist, &b->olist_hp, &b->clist, s->opt->bw_thres,
s->opt->max_n_chain, 1, &(b->k_flag), &b->r_buf, &(b->tmp_region), NULL, &(b->sp), 1, NULL);
clear_Cigar_record(&b->cigar1);
clear_Round2_alignment(&b->round2);
// return;
// b->num_correct_base += overlap_statistics(&b->olist, NULL, 0);
ha_ovec_buf_t *b = s->hab[tid];
glchain_t *bl = &(s->ll[tid]);
int64_t /**rid = s->id+i,**/ winLen = MIN((((double)THRESHOLD_MAX_SIZE)/s->opt->diff_ec_ul), WINDOW);
// uint64_t align = 0;
int fully_cov, abnormal;
assert(UL_INF.a[s->id+i].rlen == s->len[i]);
// void *km = s->buf?(s->buf[tid]?s->buf[tid]->km:NULL):NULL;
// if(s->id+i!=873) return;
// fprintf(stderr, "\n[M::%s] rid:%ld, len:%lu\n", __func__, s->id+i, s->len[i]);
// if (memcmp(UL_INF.nid.a[s->id+i].a, "d0aab024-b3a7-40fb-83cc-22c3d6d951f8", UL_INF.nid.a[s->id+i].n-1)) return;
// fprintf(stderr, "[M::%s::] ==> len: %lu\n", __func__, s->len[i]);
ha_get_ul_candidates_interface(b->abl, i, s->seq[i], s->len[i], s->opt->w, s->opt->k, s->uu, &b->olist, &b->olist_hp, &b->clist, s->opt->bw_thres,
s->opt->max_n_chain, 1, NULL/**&(b->k_flag)**/, &b->r_buf, &(b->tmp_region), NULL, &(b->sp), 1, NULL);
clear_Cigar_record(&b->cigar1);
clear_Round2_alignment(&b->round2);
// return;
// b->num_correct_base += overlap_statistics(&b->olist, NULL, 0);
b->self_read.seq = s->seq[i]; b->self_read.length = s->len[i]; b->self_read.size = 0;
correct_ul_overlap(&b->olist, s->uu, &b->self_read, &b->correct, &b->ovlp_read, &b->POA_Graph, &b->DAGCon,
&b->cigar1, &b->hap, &b->round2, 0, 1, &fully_cov, &abnormal, s->opt->diff_ec_ul, winLen, NULL);
b->self_read.seq = s->seq[i]; b->self_read.length = s->len[i]; b->self_read.size = 0;
correct_ul_overlap(&b->olist, s->uu, &b->self_read, &b->correct, &b->ovlp_read, &b->POA_Graph, &b->DAGCon,
&b->cigar1, &b->hap, &b->round2, 0, 1, &fully_cov, &abnormal, s->opt->diff_ec_ul, winLen, NULL);
// uint64_t k;
// for (k = 0; k < b->olist.length; k++) {
// if(b->olist.list[k].is_match == 1) b->num_correct_base += b->olist.list[k].x_pos_e+1-b->olist.list[k].x_pos_s;
// if(b->olist.list[k].is_match == 2) b->num_recorrect_base += b->olist.list[k].x_pos_e+1-b->olist.list[k].x_pos_s;
// }
// uint64_t k;
// for (k = 0; k < b->olist.length; k++) {
// if(b->olist.list[k].is_match == 1) b->num_correct_base += b->olist.list[k].x_pos_e+1-b->olist.list[k].x_pos_s;
// if(b->olist.list[k].is_match == 2) b->num_recorrect_base += b->olist.list[k].x_pos_e+1-b->olist.list[k].x_pos_s;
// }
// gl_chain_refine(&b->olist, &b->correct, &b->hap, bl, s->uu, s->opt->diff_ec_ul, winLen, s->len[i], km);
gl_chain_refine_advance_combine(s->buf[tid], &(UL_INF.a[i]), &b->olist, &b->correct, &b->hap, &(s->sps[tid]), bl, &(s->gdp[tid]), s->uu, s->opt->diff_ec_ul, winLen, s->len[i], s->uopt, s->id+i, tid, NULL);
// return;
// b->num_read_base += b->self_read.length;
// b->num_correct_base += b->correct.corrected_base;
// b->num_recorrect_base += b->round2.dumy.corrected_base;
memset(&b->self_read, 0, sizeof(b->self_read));
if(UL_INF.a[i].dd) {
free(s->seq[i]); s->seq[i] = NULL; b->num_correct_base++;
}
s->hab[tid]->num_read_base++;
// gl_chain_refine(&b->olist, &b->correct, &b->hap, bl, s->uu, s->opt->diff_ec_ul, winLen, s->len[i], km);
gl_chain_refine_advance_combine(s->buf[tid], &(UL_INF.a[s->id+i]), &b->olist, &b->correct, &b->hap, &(s->sps[tid]), bl, &(s->gdp[tid]), s->uu, s->opt->diff_ec_ul, winLen, s->len[i], s->uopt, s->id+i, tid, NULL);
// return;
// b->num_read_base += b->self_read.length;
// b->num_correct_base += b->correct.corrected_base;
// b->num_recorrect_base += b->round2.dumy.corrected_base;
memset(&b->self_read, 0, sizeof(b->self_read));
if(UL_INF.a[s->id+i].dd) {
free(s->seq[i]); s->seq[i] = NULL; b->num_correct_base++;
}
s->hab[tid]->num_read_base++;
// align = kv_ul_ov_t_statistics(&(bl->tk), i, &(b->num_recorrect_base));
// if(align == s->len[i]) {
// free(s->seq[i]); s->seq[i] = NULL;
// }
// b->num_correct_base += align;
int64_t mem[6], mem_hab[6];
if(get_utepdat_t_mem_tid(s, tid, mem, mem_hab)>((int64_t)5*(int64_t)1073741824)) {
fprintf(stderr, "[M::%s::tid->%d::rid->%ld] buffer[0]: %.3fGB(%.3fGB::%.3fGB::%.3fGB::%.3fGB::%.3fGB), buffer[1]: %.3fGB, buffer[2]: %.3fGB, buffer[3]: %.3fGB, buffer[4]: %.3fGB, buffer[5]: %.3fGB\n",
__func__, tid, i, mem[0]/1073741824.0,
mem_hab[0]/1073741824.0, mem_hab[1]/1073741824.0, mem_hab[2]/1073741824.0,
mem_hab[3]/1073741824.0, mem_hab[4]/1073741824.0,
mem[1]/1073741824.0, mem[2]/1073741824.0,
mem[3]/1073741824.0, mem[4]/1073741824.0, mem[5]/1073741824.0);
}
// uint64_t k;
// b->num_read_base += overlap_statistics(&b->olist, NULL, NULL, 1);
// for (k = 0; k < bl->tk.n; k++) {
// if(bl->tk.a[k].sec == 0) b->num_correct_base += bl->tk.a[k].qe - bl->tk.a[k].qs;
// if(bl->tk.a[k].sec > 0) b->num_recorrect_base += bl->tk.a[k].qe - bl->tk.a[k].qs;
// }
// for (k = 0; k < bl->lo.n; k++) {
// b->num_read_base += bl->lo.a[k].qe - bl->lo.a[k].qs;
// }
// align = kv_ul_ov_t_statistics(&(bl->tk), i, &(b->num_recorrect_base));
// if(align == s->len[i]) {
// free(s->seq[i]); s->seq[i] = NULL;
// }
// b->num_correct_base += align;
// uint32_t l1 = overlap_statistics(&b->olist, s->uu->ug, 1), l2 = overlap_statistics(&b->olist, s->uu->ug, 2);
// uint64_t k;
// b->num_read_base += overlap_statistics(&b->olist, NULL, NULL, 1);
// for (k = 0; k < bl->tk.n; k++) {
// if(bl->tk.a[k].sec == 0) b->num_correct_base += bl->tk.a[k].qe - bl->tk.a[k].qs;
// if(bl->tk.a[k].sec > 0) b->num_recorrect_base += bl->tk.a[k].qe - bl->tk.a[k].qs;
// }
// for (k = 0; k < bl->lo.n; k++) {
// b->num_read_base += bl->lo.a[k].qe - bl->lo.a[k].qs;
// }
// uint32_t l1 = overlap_statistics(&b->olist, s->uu->ug, 1), l2 = overlap_statistics(&b->olist, s->uu->ug, 2);
//
// if(l1 == 0 && l2 > 0) fprintf(stderr, "[M::%s::%lu::no_match]\n", UL_INF.nid.a[s->id+i].a, s->len[i]);
// fprintf(stderr, "[M::%s::%lu::] l1->%u; l2->%u\n", UL_INF.nid.a[s->id+i].a, s->len[i], l1, l2);
// fprintf(stderr, "[M::%s::rid->%ld] done\n", __func__, s->id+i);
// if(l1 == 0 && l2 > 0) fprintf(stderr, "[M::%s::%lu::no_match]\n", UL_INF.nid.a[s->id+i].a, s->len[i]);
// fprintf(stderr, "[M::%s::%lu::] l1->%u; l2->%u\n", UL_INF.nid.a[s->id+i].a, s->len[i], l1, l2);
// fprintf(stderr, "[M::%s::rid->%ld] done\n", __func__, s->id+i);
}
void dump_gaf(mg_gres_a *hits, const mg_gchains_t *gs, uint32_t only_p)
{
if (gs == NULL || gs->n_gc == 0 || gs->n_lc == 0) return;
@@ -5262,6 +5305,27 @@ void dump_gaf(mg_gres_a *hits, const mg_gchains_t *gs, uint32_t only_p)
}
}
int64_t get_utepdat_t_mem(const utepdat_t *b, int64_t is_print)
{
int64_t i, mem[7] = {0}, tt[6] = {0};
for (i = 0; i < b->n_thread; i++) {
get_utepdat_t_mem_tid(b, i, tt, NULL);
mem[0] += tt[0]; mem[1] += tt[1]; mem[2] += tt[2]; mem[3] += tt[3]; mem[4] += tt[4]; mem[5] += tt[5];
}
for (i = 0; i < b->n; ++i) mem[6] += b->len[i];
mem[6] += (sizeof(*(b->len))*b->n) + (sizeof(*(b->seq))*b->n);
if(is_print) {
for (i = 0; i < 7; i++) {
fprintf(stderr, "[M::%s] size of buffer[%ld]: %.3fGB, %.3fKB, %ldB\n",
__func__, i, mem[i]/1073741824.0, mem[i]/1048576.0, mem[i]);
}
}
return mem[0] + mem[1] + mem[2] + mem[3] + mem[4] + mem[5] + mem[6];
}
static void *worker_ul_pipeline(void *data, int step, void *in) // callback for kt_pipeline()
{
uldat_t *p = (uldat_t*)data;
@@ -5536,7 +5600,7 @@ static void *worker_ul_rescall_pipeline(void *data, int step, void *in) // callb
else if (step == 1) { // step 2: alignment
utepdat_t *s = (utepdat_t*)in;
uint64_t i;
uint64_t i; s->n_thread = p->n_thread;
CALLOC(s->hab, p->n_thread); CALLOC(s->ll, p->n_thread); CALLOC(s->buf, p->n_thread);
CALLOC(s->gdp, p->n_thread); CALLOC(s->mzs, p->n_thread); CALLOC(s->sps, p->n_thread);
@@ -5547,9 +5611,12 @@ static void *worker_ul_rescall_pipeline(void *data, int step, void *in) // callb
fprintf(stderr, "[M::%s::Start] ==> s->id: %lu, s->n:% d\n", __func__, s->id, s->n);
kt_for(p->n_thread, worker_for_ul_rescall_alignment, s, s->n);
fprintf(stderr, "[M::%s::Done] ==> s->id: %lu, s->n:% d\n", __func__, s->id, s->n);
get_utepdat_t_mem(s, 1);
for (i = 0; i < p->n_thread; ++i) {
p->num_bases += s->hab[i]->num_read_base;
p->num_corrected_bases += s->hab[i]->num_correct_base;
// s->num_recorrected_bases += s->hab[i]->num_recorrect_base;
ha_ovec_destroy(s->hab[i]); hc_glchain_destroy(&(s->ll[i]));
mg_tbuf_destroy(s->buf[i]); hc_gdpchain_destroy(&(s->gdp[i]));
@@ -5564,6 +5631,8 @@ static void *worker_ul_rescall_pipeline(void *data, int step, void *in) // callb
for (i = 0; i < s->n; ++i) {
rid = s->id + i;
if(UL_INF.a[rid].dd == 0 && p->ucr_s && p->ucr_s->flag == 1) {
assert(s->seq[i]);
// if(s->seq[i] == NULL) fprintf(stderr, "[M::%s::]rid->%ld, len->%lu\n", __func__, rid, s->len[i]);
write_compress_base_disk(p->ucr_s->fp, rid, s->seq[i], s->len[i], &(p->ucr_s->u));
}
free(s->seq[i]);
@@ -8540,11 +8609,12 @@ int scall_ul_pipeline(uldat_t* sl, const enzyme *fn)
int rescall_ul_pipeline(uldat_t* sl, const enzyme *fn)
{
double index_time = yak_realtime();
int32_t i; uint32_t k;
int32_t i; uint32_t k, rlen;
///clean UL_INF
for (k = 0; k < UL_INF.n; k++) {
rlen = UL_INF.a[k].rlen;
free(UL_INF.a[k].bb.a); free(UL_INF.a[k].N_site.a); free(UL_INF.a[k].r_base.a);
memset(&(UL_INF.a[k]), 0, sizeof(UL_INF.a[k]));
memset(&(UL_INF.a[k]), 0, sizeof(UL_INF.a[k])); UL_INF.a[k].rlen = rlen;
}
for (i = 0; i < fn->n; i++){