diff --git a/oonipipeline/src/oonipipeline/tasks/updaters/asnmeta_updater.py b/oonipipeline/src/oonipipeline/tasks/updaters/asnmeta_updater.py index 1822f402c..cface62ae 100644 --- a/oonipipeline/src/oonipipeline/tasks/updaters/asnmeta_updater.py +++ b/oonipipeline/src/oonipipeline/tasks/updaters/asnmeta_updater.py @@ -17,6 +17,27 @@ AS_ORG_MAP_URL = "https://archive.org/download/ip2country-as/all_as_org_map.json" +# asnmeta lives on the replicated oonidata_cluster, as a permanent pair of +# tables -- asnmeta and asnmeta_tmp -- matching the same swap-table pattern +# used for citizenlab/citizenlab_flip. This script creates both tables +# itself (idempotently) if they don't already exist. +# +# TRUNCATE and INSERT are data operations: ClickHouse replicates those +# automatically to every replica via the table's own replication log, no +# ON CLUSTER needed. CREATE and EXCHANGE TABLES are different -- they're +# catalog-level operations, and `ooni` is a plain Atomic database, so table +# *names* are local to each node's own catalog and do not follow the +# table's data replication. Without ON CLUSTER on these, the swap would +# only rename things on whichever single node this script's client +# connects to -- the other replicas would keep calling the OLD data +# "asnmeta" indefinitely, every single run. +# +# Plain ReplicatedMergeTree, not Replacing: asnmeta intentionally keeps +# every historical row per ASN (queries pick the latest via +# `changed`/argMax at read time) rather than relying on background merges +# to dedup them away. +CLUSTER_NAME = "oonidata_cluster" + log = logging.getLogger("analysis.asnmeta_updater") # metrics = setup_metrics(name="asnmeta_updater") progress_cnt = 0 @@ -57,35 +78,31 @@ def fetch_data() -> List[dict]: def update_asnmeta(clickhouse_url: str) -> None: progress("starting") click = Clickhouse.from_url(clickhouse_url) - q = """ - CREATE TABLE IF NOT EXISTS asnmeta ( - asn UInt32, - org_name String, - cc String, - changed Date, - aut_name String, - source String - ) ENGINE = MergeTree() - ORDER BY (asn, changed) - """ - click.execute(q) - q = "DROP TABLE IF EXISTS asnmeta_tmp" - click.execute(q) - - q = """ - CREATE TABLE asnmeta_tmp ( - asn UInt32, - org_name String, - cc String, - changed Date, - aut_name String, - source String - ) ENGINE = MergeTree() - ORDER BY (asn, changed) - """ + # Ensure both tables exist (first-ever run, a fresh test/dev + # environment, or after manual recovery). IF NOT EXISTS makes this a + # no-op on every normal run once they're already there -- unlike a + # DROP+CREATE-every-run pattern, this doesn't churn each table's + # ZooKeeper replication metadata on every single scheduled update. + for table_name in ("asnmeta", "asnmeta_tmp"): + q = f""" + CREATE TABLE IF NOT EXISTS {table_name} ON CLUSTER {CLUSTER_NAME} ( + asn UInt32, + org_name String, + cc String, + changed Date, + aut_name String, + source String + ) ENGINE = ReplicatedMergeTree('/clickhouse/{{cluster}}/tables/ooni/{table_name}/{{shard}}', '{{replica}}') + ORDER BY (asn, changed) + """ + click.execute(q) + progress("asnmeta/asnmeta_tmp ensured") + + log.info("Emptying Clickhouse asnmeta_tmp table") + q = "TRUNCATE TABLE asnmeta_tmp" click.execute(q) - progress("asnmeta_tmp recreated") + progress("asnmeta_tmp truncated") log.info(f"Ingesting {AS_ORG_MAP_URL}") data = fetch_data() @@ -106,6 +123,6 @@ def update_asnmeta(clickhouse_url: str) -> None: assert 100_000 < row_cnt < 1_000_000 log.info("Swapping tables") - q = "EXCHANGE TABLES asnmeta_tmp AND asnmeta" + q = f"EXCHANGE TABLES asnmeta_tmp AND asnmeta ON CLUSTER {CLUSTER_NAME}" click.execute(q) progress("asnmeta ready") diff --git a/oonipipeline/tests/clickhouse-config.d/cluster.xml b/oonipipeline/tests/clickhouse-config.d/cluster.xml new file mode 100644 index 000000000..44d459504 --- /dev/null +++ b/oonipipeline/tests/clickhouse-config.d/cluster.xml @@ -0,0 +1,55 @@ + + + + + + + localhost + 9000 + + + + + + + oonidata_cluster + 01 + 01 + + + + 9181 + 1 + /var/lib/clickhouse/coordination/log + /var/lib/clickhouse/coordination/snapshots + + 10000 + 30000 + + + + 1 + localhost + 9234 + + + + + + + localhost + 9181 + + + + + /clickhouse/task_queue/ddl + + diff --git a/oonipipeline/tests/docker-compose.yml b/oonipipeline/tests/docker-compose.yml index 220c74eb7..68dbb2880 100644 --- a/oonipipeline/tests/docker-compose.yml +++ b/oonipipeline/tests/docker-compose.yml @@ -6,4 +6,12 @@ services: - "9000" environment: CLICKHOUSE_PASSWORD: test - CLICKHOUSE_USER: test \ No newline at end of file + CLICKHOUSE_USER: test + volumes: + - ./clickhouse-config.d:/etc/clickhouse-server/config.d:ro + healthcheck: + test: ["CMD", "clickhouse-client", "--query", "select 1;"] + interval: 5s + retries: 10 + start_period: 30s + timeout: 10s \ No newline at end of file diff --git a/oonipipeline/tests/test_updaters.py b/oonipipeline/tests/test_updaters.py index bd8c76ef0..5c636081c 100644 --- a/oonipipeline/tests/test_updaters.py +++ b/oonipipeline/tests/test_updaters.py @@ -57,11 +57,12 @@ def mocked_execute(q, data=None, **kw): qrs = [" ".join(q[0][0].split()) for q in mock_click.execute.call_args_list] expected_queries = [ - "DROP TABLE IF EXISTS asnmeta_tmp", - "CREATE TABLE asnmeta_tmp ( asn UInt32, org_name String, cc String, changed Date, aut_name String, source String ) ENGINE = MergeTree() ORDER BY (asn, changed)", + "CREATE TABLE IF NOT EXISTS asnmeta ON CLUSTER oonidata_cluster ( asn UInt32, org_name String, cc String, changed Date, aut_name String, source String ) ENGINE = ReplicatedMergeTree('/clickhouse/{cluster}/tables/ooni/asnmeta/{shard}', '{replica}') ORDER BY (asn, changed)", + "CREATE TABLE IF NOT EXISTS asnmeta_tmp ON CLUSTER oonidata_cluster ( asn UInt32, org_name String, cc String, changed Date, aut_name String, source String ) ENGINE = ReplicatedMergeTree('/clickhouse/{cluster}/tables/ooni/asnmeta_tmp/{shard}', '{replica}') ORDER BY (asn, changed)", + "TRUNCATE TABLE asnmeta_tmp", "INSERT INTO asnmeta_tmp (asn, org_name, cc, changed, aut_name, source) VALUES", "SELECT count() FROM asnmeta_tmp", - "EXCHANGE TABLES asnmeta_tmp AND asnmeta", + "EXCHANGE TABLES asnmeta_tmp AND asnmeta ON CLUSTER oonidata_cluster", ] for qr in expected_queries: assert qr in qrs