From 5580f164cf7853e1c2396866415f1e6efb49da5b Mon Sep 17 00:00:00 2001 From: nitowa Date: Wed, 24 Aug 2022 11:55:09 -0400 Subject: [PATCH] rewrite logic to reduce DB reads --- src/spark/main.py | 31 ------------------------------- 1 file changed, 31 deletions(-) diff --git a/src/spark/main.py b/src/spark/main.py index 607cca1..f41f724 100644 --- a/src/spark/main.py +++ b/src/spark/main.py @@ -94,37 +94,6 @@ def find(data: tuple[Row, Iterable[str]]) -> str | None: else: 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) tx_addr_groups = master.group_tx_addrs()