@@ -69,9 +69,9 @@ namespace graphchi {
6969 std::string blockfilename;
7070 dynamicdata_block<ET > * dynblock;
7171
72- sblock () : writedesc(0 ), readdesc(0 ), active(false ) { data = NULL ; }
72+ sblock () : writedesc(0 ), readdesc(0 ), active(false ) { data = NULL ; dynblock = NULL ; }
7373 sblock (int wdesc, int rdesc, bool is_edata_block=false ) : writedesc(wdesc), readdesc(rdesc), active(false ),
74- is_edata_block (is_edata_block){ data = NULL ; }
74+ is_edata_block (is_edata_block){ data = NULL ; dynblock = NULL ; }
7575 sblock (int wdesc, int rdesc, bool is_edata_block, std::string blockfilename) : writedesc(wdesc), readdesc(rdesc), active(false ),
7676 is_edata_block(is_edata_block), blockfilename(blockfilename) {
7777 assert (is_edata_block == true );
@@ -90,7 +90,8 @@ namespace graphchi {
9090 int realsize;
9191 dynblock->write (&outdata, realsize);
9292 write_block_uncompressed_size (blockfilename, realsize);
93- iomgr->managed_pwritea_now (writedesc, &data, realsize, 0 ); /* Need to write whole block in the compressed regime */
93+ iomgr->managed_pwritea_now (writedesc, &outdata, realsize, 0 ); /* Need to write whole block in the compressed regime */
94+ free (outdata);
9495 } else {
9596 iomgr->managed_pwritea_now (writedesc, &data, len, offset);
9697 }
@@ -104,7 +105,7 @@ namespace graphchi {
104105 int realsize = get_block_uncompressed_size (blockfilename, end-offset);
105106 iomgr->managed_preada_now (readdesc, &data, realsize, 0 );
106107 int nedges = (end - offset) / sizeof (int ); // Ugly
107- dynblock = new dynamicdata_block<ET >(nedges, (uint8_t *) & data);
108+ dynblock = new dynamicdata_block<ET >(nedges, (uint8_t *) data, realsize );
108109 } else {
109110 iomgr->managed_preada_now (readdesc, &data, end - offset, offset);
110111 }
@@ -280,10 +281,13 @@ namespace graphchi {
280281 size_t correction = edataoffset - newblock.offset ;
281282 newblock.end = std::min (edatafilesize, newblock.offset + blocksize);
282283 assert (newblock.end >= newblock.offset );
283- iomgr->managed_malloc (edata_session, &newblock.data , newblock.end - newblock.offset , newblock.offset );
284+ int realsize = get_block_uncompressed_size (blockfilename, newblock.end - newblock.offset );
285+ iomgr->managed_malloc (edata_session, &newblock.data , realsize, newblock.offset );
284286 newblock.ptr = newblock.data + correction;
285287 activeblocks.push_back (newblock);
286288 curblock = &activeblocks[activeblocks.size ()-1 ];
289+ curblock->active = true ;
290+ curblock->read_now (iomgr);
287291 }
288292 }
289293
@@ -321,8 +325,10 @@ namespace graphchi {
321325 if (only_adjacency) return NULL ;
322326 check_curblock (sizeof (int ));
323327 edataoffset += sizeof (int );
328+ int blockedgeidx = (curblock->ptr - curblock->data ) / sizeof (int );
324329 curblock->ptr += sizeof (int );
325- return curblock->dynblock ->edgevec ((curblock->ptr - curblock->data ) / sizeof (int ));
330+ assert (curblock->dynblock != NULL );
331+ return curblock->dynblock ->edgevec (blockedgeidx);
326332 }
327333
328334 inline void skip (int n, int sz) {
@@ -394,18 +400,8 @@ namespace graphchi {
394400 bool special_edge = false ;
395401 vid_t target = (sizeof (ET ) == sizeof (ETspecial) ? read_val<vid_t >() : translate_edge (read_val<vid_t >(), special_edge));
396402 ET * evalue = read_edgeptr ();
403+
397404
398- if (!only_adjacency) {
399- if (!curblock->active ) {
400- if (async_edata_loading) {
401- curblock->read_async (iomgr);
402- } else {
403- curblock->read_now (iomgr);
404- }
405- }
406- // Note: this needs to be set always because curblock might change during this loop.
407- curblock->active = true ; // This block has an scheduled vertex - need to commit
408- }
409405 vertex.add_outedge (target, evalue, special_edge);
410406
411407 if (!((target >= range_st && target <= range_end))) {
0 commit comments