Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
73 changes: 45 additions & 28 deletions oonipipeline/src/oonipipeline/tasks/updaters/asnmeta_updater.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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()
Expand All @@ -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")
55 changes: 55 additions & 0 deletions oonipipeline/tests/clickhouse-config.d/cluster.xml
Original file line number Diff line number Diff line change
@@ -0,0 +1,55 @@
<!--
Single-node stand-in for the real 3-replica oonidata_cluster (production:
clickhouse1/2/3.prod.ooni.io). One shard, one replica, pointing at itself,
with an embedded Keeper for coordination: enough for `ON CLUSTER
oonidata_cluster` DDL and `ReplicatedMergeTree`/`{cluster}`/{shard}`/
`{replica}` macros to behave the same way they do against the real
cluster, without needing multiple containers in CI.
-->
<clickhouse>
<remote_servers>
<oonidata_cluster>
<shard>
<replica>
<host>localhost</host>
<port>9000</port>
</replica>
</shard>
</oonidata_cluster>
</remote_servers>

<macros>
<cluster>oonidata_cluster</cluster>
<shard>01</shard>
<replica>01</replica>
</macros>

<keeper_server>
<tcp_port>9181</tcp_port>
<server_id>1</server_id>
<log_storage_path>/var/lib/clickhouse/coordination/log</log_storage_path>
<snapshot_storage_path>/var/lib/clickhouse/coordination/snapshots</snapshot_storage_path>
<coordination_settings>
<operation_timeout_ms>10000</operation_timeout_ms>
<session_timeout_ms>30000</session_timeout_ms>
</coordination_settings>
<raft_configuration>
<server>
<id>1</id>
<hostname>localhost</hostname>
<port>9234</port>
</server>
</raft_configuration>
</keeper_server>

<zookeeper>
<node>
<host>localhost</host>
<port>9181</port>
</node>
</zookeeper>

<distributed_ddl>
<path>/clickhouse/task_queue/ddl</path>
</distributed_ddl>
</clickhouse>
10 changes: 9 additions & 1 deletion oonipipeline/tests/docker-compose.yml
Original file line number Diff line number Diff line change
Expand Up @@ -6,4 +6,12 @@ services:
- "9000"
environment:
CLICKHOUSE_PASSWORD: test
CLICKHOUSE_USER: test
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
7 changes: 4 additions & 3 deletions oonipipeline/tests/test_updaters.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading