rewrite logic to reduce DB reads
This commit is contained in:
@@ -94,37 +94,6 @@ def find(data: tuple[Row, Iterable[str]]) -> str | None:
|
|||||||
else:
|
else:
|
||||||
return None
|
return None
|
||||||
|
|
||||||
def handleTx(tx_addr_group: Row):
|
|
||||||
|
|
||||||
|
|
||||||
found_clusters: "RDD[str]" = clusters.rdd \
|
|
||||||
.map(lambda cluster: (cluster, tx_addr_group['addresses'])) \
|
|
||||||
.map(find) \
|
|
||||||
.filter(lambda x: x != None)
|
|
||||||
|
|
||||||
|
|
||||||
if(found_clusters.count() == 0):
|
|
||||||
insertNewCluster(tx_addr_group)
|
|
||||||
return
|
|
||||||
|
|
||||||
cluster_roots = found_clusters.collect()
|
|
||||||
|
|
||||||
cl = clusters \
|
|
||||||
.select('addresses') \
|
|
||||||
.where(
|
|
||||||
F.col('parent').isin(cluster_roots)
|
|
||||||
) \
|
|
||||||
.agg(F.collect_set('addresses').alias('agg')) \
|
|
||||||
.select(F.flatten('agg').alias('addresses')) \
|
|
||||||
.select(F.explode('addresses')) \
|
|
||||||
.rdd \
|
|
||||||
.map(lambda addr: (addr, cluster_roots[0])) \
|
|
||||||
.toDF(['address', 'parent']) \
|
|
||||||
.show()
|
|
||||||
#.writeTo(CLUSTERS_TABLE) \
|
|
||||||
#.append()
|
|
||||||
|
|
||||||
|
|
||||||
master = Master(config)
|
master = Master(config)
|
||||||
|
|
||||||
tx_addr_groups = master.group_tx_addrs()
|
tx_addr_groups = master.group_tx_addrs()
|
||||||
|
|||||||
Reference in New Issue
Block a user