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()