Skip to content

Commit fe83df4

Browse files
author
Aapo Kyrola
committed
working dynamic edgedata? Still a bit of a hack though....
--HG-- branch : research
1 parent e731da3 commit fe83df4

4 files changed

Lines changed: 33 additions & 27 deletions

File tree

‎src/shards/dynamicdata/dynamicblock.hpp‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -64,10 +64,11 @@ namespace graphchi {
6464

6565
dynamicdata_block() : data(NULL), chivecs(NULL) {}
6666

67-
dynamicdata_block(int nedges, uint8_t * data) {
67+
dynamicdata_block(int nedges, uint8_t * data, int datasize) : nedges(nedges){
6868
chivecs = new ET[nedges];
6969
uint8_t * ptr = data;
7070
for(int i=0; i < nedges; i++) {
71+
assert(ptr - data <= datasize);
7172
uint16_t * sz = ((uint16_t *) ptr);
7273
ptr += sizeof(uint16_t);
7374
chivecs[i] = ET(sz, (typename ET::element_type_t *) ptr);
@@ -76,6 +77,8 @@ namespace graphchi {
7677
}
7778

7879
ET * edgevec(int i) {
80+
assert(i < nedges);
81+
assert(chivecs != NULL);
7982
return &chivecs[i];
8083
}
8184

‎src/shards/dynamicdata/memoryshard.hpp‎

Lines changed: 8 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -70,7 +70,7 @@ namespace graphchi {
7070
uint8_t * adjdata;
7171
char ** edgedata;
7272
std::vector<size_t> blocksizes;
73-
std::vector< dynamicdata_block<ET> > dynamicblocks;
73+
std::vector< dynamicdata_block<ET> * > dynamicblocks;
7474
uint64_t chunkid;
7575

7676
std::vector<int> block_edatasessions;
@@ -107,8 +107,10 @@ namespace graphchi {
107107
if (edgedata[i] != NULL) {
108108
iomgr->managed_release(block_edatasessions[i], &edgedata[i]);
109109
iomgr->close_session(block_edatasessions[i]);
110+
delete dynamicblocks[i];
110111
}
111112
}
113+
dynamicblocks.clear();
112114
if (adj_session >= 0) {
113115
if (adjdata != NULL) iomgr->managed_release(adj_session, &adjdata);
114116
iomgr->close_session(adj_session);
@@ -121,10 +123,10 @@ namespace graphchi {
121123
void write_and_release_block(int i) {
122124
std::string block_filename = filename_shard_edata_block(filename_edata, i, blocksize);
123125

124-
dynamicdata_block<ET> & dynblock = dynamicblocks[i];
126+
dynamicdata_block<ET> * dynblock = dynamicblocks[i];
125127
uint8_t * outdata;
126128
int outsize;
127-
dynblock.write(&outdata, outsize);
129+
dynblock->write(&outdata, outsize);
128130
write_block_uncompressed_size(block_filename, outsize);
129131
iomgr->managed_pwritea_now(block_edatasessions[i], &outdata, outsize, 0);
130132
iomgr->managed_release(block_edatasessions[i], &edgedata[i]);
@@ -220,7 +222,7 @@ namespace graphchi {
220222
} else {
221223
iomgr->managed_preada_now(blocksession, &edgedata[blockid], fsize, 0);
222224
}
223-
dynamicblocks.push_back(dynamicdata_block<ET>(nedges, (uint8_t*) edgedata[blockid]));
225+
dynamicblocks.push_back(new dynamicdata_block<ET>(nedges, (uint8_t*) edgedata[blockid], fsize));
224226

225227
blockid++;
226228

@@ -349,7 +351,7 @@ namespace graphchi {
349351
ptr += sizeof(vid_t);
350352
if (vertex != NULL && outedges)
351353
{
352-
vertex->add_outedge(target, (only_adjacency ? NULL : dynamicblocks[blockid].edgevec(edgeptr % blocksize)), false);
354+
vertex->add_outedge(target, (only_adjacency ? NULL : dynamicblocks[blockid]->edgevec((edgeptr % blocksize)/sizeof(int))), false);
353355
}
354356

355357
if (target >= window_st) {
@@ -359,7 +361,7 @@ namespace graphchi {
359361
if (dstvertex.scheduled) {
360362
any_edges = true;
361363
// assert(only_adjacency || edgeptr < edatafilesize);
362-
ET * eptr = (only_adjacency ? NULL : dynamicblocks[blockid].edgevec(edgeptr % blocksize));
364+
ET * eptr = (only_adjacency ? NULL : dynamicblocks[blockid]->edgevec((edgeptr % blocksize)/sizeof(int)));
363365

364366
dstvertex.add_inedge(vid, (only_adjacency ? NULL : eptr), false);
365367
dstvertex.parallel_safe = dstvertex.parallel_safe && (vertex == NULL); // Avoid if

‎src/shards/dynamicdata/slidingshard.hpp‎

Lines changed: 13 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -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))) {

‎src/tests/dynamicdata_smoketest.cpp‎

Lines changed: 8 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -63,16 +63,21 @@ struct DynamicDataSmokeTestProgram : public GraphChiProgram<VertexDataType, Edge
6363

6464
evector->add(vertex.id());
6565
assert(evector->size() == 1);
66-
66+
assert(evector->get(0) == vertex.id());
6767
}
6868

6969
} else {
7070
for(int i=0; i < vertex.num_inedges(); i++) {
7171
graphchi_edge<EdgeDataType> * edge = vertex.inedge(i);
72-
chivector<vid_t> * evector = vertex.outedge(i)->get_vector();
72+
chivector<vid_t> * evector = edge->get_vector();
7373
assert(evector->size() >= gcontext.iteration);
7474
for(int j=0; j < evector->size(); j++) {
75-
assert(evector->get(j) == edge->vertex_id() + j);
75+
vid_t expected =edge->vertex_id() + j;
76+
vid_t has = evector->get(j);
77+
if (has != expected) {
78+
std::cout << "Mismatch: " << has << " != " << expected << std::endl;
79+
}
80+
assert(has == expected);
7681
}
7782
}
7883
for(int i=0; i < vertex.num_outedges(); i++) {

0 commit comments

Comments
 (0)