From 84ca8827c51f895b3e41b6528b98203be46325f9 Mon Sep 17 00:00:00 2001 From: Artur Pata Date: Mon, 27 Jul 2026 15:22:06 +0300 Subject: [PATCH 1/2] Fix NOT_ENOUGH_SPACE error from CH during monthly clean sites job --- lib/workers/clickhouse_clean_sites.ex | 81 +++++++++++++++++--- test/workers/clickhouse_clean_sites_test.exs | 66 ++++++++++++++-- 2 files changed, 130 insertions(+), 17 deletions(-) diff --git a/lib/workers/clickhouse_clean_sites.ex b/lib/workers/clickhouse_clean_sites.ex index a425fdfa5864..2332a2f5177f 100644 --- a/lib/workers/clickhouse_clean_sites.ex +++ b/lib/workers/clickhouse_clean_sites.ex @@ -1,9 +1,18 @@ defmodule Plausible.Workers.ClickhouseCleanSites do @moduledoc """ - Cleans deleted site data from ClickHouse asynchronously. + Cleans deleted site data from ClickHouse sequentially. - We batch up data deletions from ClickHouse as deleting a single site is - just as expensive as deleting many. + We batch up data deletions from ClickHouse as deleting a single site may be more expensive than deleting many: + the expense scales with the number of partitions the site or sites have rows in. + A single site with events since 2025-01 to 2026-06 will be more expensive to clean up + than a hundred sites with events only in the partition 2026-06. + + Tables are cleaned one partition at a time because to clean one partition, + Clickhouse needs to rewrite it. It reserves room for rewriting it in full on disk. + Cleaning all partitions at the same time would mean Clickhouse reserving room on disk + equal to the size of all the partitions, + potentially reserving all the available disk space and shorting out INSERTs from ingestion. + This sequential approach prevents that. """ use Plausible.Repo @@ -31,7 +40,9 @@ defmodule Plausible.Workers.ClickhouseCleanSites do "imported_visitors" ] - @settings if Mix.env() in [:test, :ce_test, :e2e_test], do: [mutations_sync: 2], else: [] + @settings [mutations_sync: 2] + + @partition_delete_timeout :timer.minutes(15) def perform(_job) do deleted_sites = get_deleted_sites_with_clickhouse_data() @@ -41,18 +52,59 @@ defmodule Plausible.Workers.ClickhouseCleanSites do "Clearing ClickHouse data for the following #{length(deleted_sites)} sites which have been deleted: #{inspect(deleted_sites)}" ) + database = current_database() + for table <- @tables_to_clear do - IngestRepo.query!( - "ALTER TABLE {$0:Identifier} DELETE WHERE site_id IN {$1:Array(UInt64)}", - [table, deleted_sites], - settings: @settings - ) + clean_sites_from_table(database, table, deleted_sites) end end :ok end + defp clean_sites_from_table(database, table, deleted_sites_ids) do + for partition_id <- active_partition_ids(database, table) do + delete_sites_from_partition(table, partition_id, deleted_sites_ids) + end + end + + defp delete_sites_from_partition(table, partition_id, deleted_sites_ids) do + IngestRepo.query!( + "ALTER TABLE {$0:Identifier} DELETE IN PARTITION ID {$1:String} WHERE site_id IN {$2:Array(UInt64)}", + [table, partition_id, deleted_sites_ids], + settings: @settings, + timeout: @partition_delete_timeout, + checkout_retries: 0 + ) + end + + defp active_partition_ids(database, table) do + source = + if IngestRepo.clustered_table?(table) do + "clusterAllReplicas('{cluster}', system.parts)" + else + "system.parts" + end + + %Ch.Result{columns: ["partition_id"], rows: rows} = + IngestRepo.query!( + """ + SELECT DISTINCT partition_id + FROM #{source} + WHERE database = {$0:String} AND table = {$1:String} AND active + ORDER BY partition_id + """, + [database, table] + ) + + Enum.map(rows, fn [partition_id] -> partition_id end) + end + + defp current_database do + %Ch.Result{rows: [[database]]} = IngestRepo.query!("SELECT currentDatabase()") + database + end + def get_deleted_sites_with_clickhouse_data() do pg_sites = from(s in Plausible.Site.regular(), select: s.id) @@ -67,15 +119,22 @@ defmodule Plausible.Workers.ClickhouseCleanSites do DBConnection.run( ch, fn conn -> - Ch.query!(conn, "FROM events_v2 SELECT site_id GROUP BY site_id", [], + Ch.query!(conn, site_ids_with_data_query(), [], + settings: [optimize_distinct_in_order: 1], timeout: :infinity ) end, timeout: :infinity ) - ch_sites = rows |> MapSet.new(fn [site_id] -> site_id end) + ch_sites = for [site_id] <- rows, not is_nil(site_id), into: MapSet.new(), do: site_id MapSet.difference(ch_sites, pg_sites) |> MapSet.to_list() end + + defp site_ids_with_data_query do + Enum.map_join(@tables_to_clear, " UNION DISTINCT ", fn table -> + "SELECT DISTINCT site_id FROM #{table}" + end) + end end diff --git a/test/workers/clickhouse_clean_sites_test.exs b/test/workers/clickhouse_clean_sites_test.exs index 18cc2c96f759..679998ae81d6 100644 --- a/test/workers/clickhouse_clean_sites_test.exs +++ b/test/workers/clickhouse_clean_sites_test.exs @@ -10,12 +10,25 @@ defmodule Plausible.Workers.ClickhouseCleanSitesTest do deleted_site = insert(:site) populate_stats(site, [ - build(:pageview) + build(:pageview), + build(:pageview, timestamp: ~D[2026-01-01]), + build(:imported_visitors), + build(:imported_sources), + build(:imported_pages), + build(:imported_entry_pages), + build(:imported_exit_pages), + build(:imported_locations), + build(:imported_devices), + build(:imported_browsers), + build(:imported_operating_systems), + build(:imported_custom_events) ]) populate_stats(deleted_site, [ build(:pageview), - build(:pageview), + build(:pageview, timestamp: ~D[2026-01-01]), + build(:pageview, timestamp: ~D[2026-02-01]), + build(:pageview, timestamp: ~D[2026-03-01]), build(:imported_visitors), build(:imported_sources), build(:imported_pages), @@ -24,7 +37,8 @@ defmodule Plausible.Workers.ClickhouseCleanSitesTest do build(:imported_locations), build(:imported_devices), build(:imported_browsers), - build(:imported_operating_systems) + build(:imported_operating_systems), + build(:imported_custom_events) ]) Repo.delete!(deleted_site) @@ -52,8 +66,19 @@ defmodule Plausible.Workers.ClickhouseCleanSitesTest do assert_count(deleted_site, "imported_devices", 0) assert_count(deleted_site, "imported_browsers", 0) assert_count(deleted_site, "imported_operating_systems", 0) - assert_count(site, "events_v2", 1) - assert_count(site, "sessions_v2", 1) + assert_count(deleted_site, "imported_custom_events", 0) + assert_count(site, "events_v2", 2) + assert_count(site, "sessions_v2", 2) + assert_count(site, "imported_visitors", 1) + assert_count(site, "imported_sources", 1) + assert_count(site, "imported_pages", 1) + assert_count(site, "imported_entry_pages", 1) + assert_count(site, "imported_exit_pages", 1) + assert_count(site, "imported_locations", 1) + assert_count(site, "imported_devices", 1) + assert_count(site, "imported_browsers", 1) + assert_count(site, "imported_operating_systems", 1) + assert_count(site, "imported_custom_events", 1) assert not Enum.member?( ClickhouseCleanSites.get_deleted_sites_with_clickhouse_data(), @@ -61,8 +86,37 @@ defmodule Plausible.Workers.ClickhouseCleanSitesTest do ) end + @tag :slow + test "cleans a deleted site that has only imported data and no native events" do + kept_site = insert(:site) + imported_only = insert(:site) + + # A live site with native events - must be left untouched. + populate_stats(kept_site, [build(:pageview)]) + + populate_stats(imported_only, [ + build(:imported_visitors), + build(:imported_sources), + build(:imported_custom_events) + ]) + + Repo.delete!(imported_only) + + assert Enum.member?( + ClickhouseCleanSites.get_deleted_sites_with_clickhouse_data(), + imported_only.id + ) + + ClickhouseCleanSites.perform(nil) + + assert_count(imported_only, "imported_visitors", 0) + assert_count(imported_only, "imported_sources", 0) + assert_count(imported_only, "imported_custom_events", 0) + assert_count(kept_site, "events_v2", 1) + end + def assert_count(site, table, expected_count) do q = from(e in table, select: %{count: fragment("count()")}, where: e.site_id == ^site.id) - await_clickhouse_count(q, expected_count) + assert await_clickhouse_count(q, expected_count) end end From fd76b386be9f031a9fc904d2ad34dbf5cae8ab01 Mon Sep 17 00:00:00 2001 From: Artur Pata Date: Thu, 30 Jul 2026 20:58:38 +0300 Subject: [PATCH 2/2] Remove job entirely --- config/runtime.exs | 3 - lib/workers/clickhouse_clean_sites.ex | 140 ------------------- test/workers/clickhouse_clean_sites_test.exs | 122 ---------------- 3 files changed, 265 deletions(-) delete mode 100644 lib/workers/clickhouse_clean_sites.ex delete mode 100644 test/workers/clickhouse_clean_sites_test.exs diff --git a/config/runtime.exs b/config/runtime.exs index 91f056353d27..99d901e080a0 100644 --- a/config/runtime.exs +++ b/config/runtime.exs @@ -835,8 +835,6 @@ cloud_cron = [ {"0 0 * * *", Plausible.Workers.LockSites}, # Daily at 8 {"0 8 * * *", Plausible.Workers.AcceptTrafficUntil}, - # First sunday of the month, 4:00 UTC - {"0 4 1-7 * SUN", Plausible.Workers.ClickhouseCleanSites}, # Daily at 4:00 UTC {"0 4 * * *", Plausible.Workers.SetLegacyTimeOnPageCutoff}, # Daily at 2:00 UTC @@ -859,7 +857,6 @@ base_queues = [ notify_exported_analytics: 1, domain_change_transition: 1, check_accept_traffic_until: 1, - clickhouse_clean_sites: 1, locations_sync: 1 ] diff --git a/lib/workers/clickhouse_clean_sites.ex b/lib/workers/clickhouse_clean_sites.ex deleted file mode 100644 index 2332a2f5177f..000000000000 --- a/lib/workers/clickhouse_clean_sites.ex +++ /dev/null @@ -1,140 +0,0 @@ -defmodule Plausible.Workers.ClickhouseCleanSites do - @moduledoc """ - Cleans deleted site data from ClickHouse sequentially. - - We batch up data deletions from ClickHouse as deleting a single site may be more expensive than deleting many: - the expense scales with the number of partitions the site or sites have rows in. - A single site with events since 2025-01 to 2026-06 will be more expensive to clean up - than a hundred sites with events only in the partition 2026-06. - - Tables are cleaned one partition at a time because to clean one partition, - Clickhouse needs to rewrite it. It reserves room for rewriting it in full on disk. - Cleaning all partitions at the same time would mean Clickhouse reserving room on disk - equal to the size of all the partitions, - potentially reserving all the available disk space and shorting out INSERTs from ingestion. - This sequential approach prevents that. - """ - - use Plausible.Repo - use Plausible.ClickhouseRepo - use Plausible.IngestRepo - use Oban.Worker, queue: :clickhouse_clean_sites - - import Ecto.Query - - require Logger - - @tables_to_clear [ - "events_v2", - "sessions_v2", - "ingest_counters", - "imported_browsers", - "imported_devices", - "imported_entry_pages", - "imported_exit_pages", - "imported_locations", - "imported_operating_systems", - "imported_pages", - "imported_custom_events", - "imported_sources", - "imported_visitors" - ] - - @settings [mutations_sync: 2] - - @partition_delete_timeout :timer.minutes(15) - - def perform(_job) do - deleted_sites = get_deleted_sites_with_clickhouse_data() - - if not Enum.empty?(deleted_sites) do - Logger.notice( - "Clearing ClickHouse data for the following #{length(deleted_sites)} sites which have been deleted: #{inspect(deleted_sites)}" - ) - - database = current_database() - - for table <- @tables_to_clear do - clean_sites_from_table(database, table, deleted_sites) - end - end - - :ok - end - - defp clean_sites_from_table(database, table, deleted_sites_ids) do - for partition_id <- active_partition_ids(database, table) do - delete_sites_from_partition(table, partition_id, deleted_sites_ids) - end - end - - defp delete_sites_from_partition(table, partition_id, deleted_sites_ids) do - IngestRepo.query!( - "ALTER TABLE {$0:Identifier} DELETE IN PARTITION ID {$1:String} WHERE site_id IN {$2:Array(UInt64)}", - [table, partition_id, deleted_sites_ids], - settings: @settings, - timeout: @partition_delete_timeout, - checkout_retries: 0 - ) - end - - defp active_partition_ids(database, table) do - source = - if IngestRepo.clustered_table?(table) do - "clusterAllReplicas('{cluster}', system.parts)" - else - "system.parts" - end - - %Ch.Result{columns: ["partition_id"], rows: rows} = - IngestRepo.query!( - """ - SELECT DISTINCT partition_id - FROM #{source} - WHERE database = {$0:String} AND table = {$1:String} AND active - ORDER BY partition_id - """, - [database, table] - ) - - Enum.map(rows, fn [partition_id] -> partition_id end) - end - - defp current_database do - %Ch.Result{rows: [[database]]} = IngestRepo.query!("SELECT currentDatabase()") - database - end - - def get_deleted_sites_with_clickhouse_data() do - pg_sites = - from(s in Plausible.Site.regular(), select: s.id) - |> Plausible.Repo.all() - |> MapSet.new() - - {:ok, ch} = - Plausible.ClickhouseRepo.get_config_without_ch_query_execution_timeout() - |> Ch.start_link() - - %Ch.Result{columns: ["site_id"], rows: rows} = - DBConnection.run( - ch, - fn conn -> - Ch.query!(conn, site_ids_with_data_query(), [], - settings: [optimize_distinct_in_order: 1], - timeout: :infinity - ) - end, - timeout: :infinity - ) - - ch_sites = for [site_id] <- rows, not is_nil(site_id), into: MapSet.new(), do: site_id - - MapSet.difference(ch_sites, pg_sites) |> MapSet.to_list() - end - - defp site_ids_with_data_query do - Enum.map_join(@tables_to_clear, " UNION DISTINCT ", fn table -> - "SELECT DISTINCT site_id FROM #{table}" - end) - end -end diff --git a/test/workers/clickhouse_clean_sites_test.exs b/test/workers/clickhouse_clean_sites_test.exs deleted file mode 100644 index 679998ae81d6..000000000000 --- a/test/workers/clickhouse_clean_sites_test.exs +++ /dev/null @@ -1,122 +0,0 @@ -defmodule Plausible.Workers.ClickhouseCleanSitesTest do - use Plausible.DataCase - import Plausible.Factory - - alias Plausible.Workers.ClickhouseCleanSites - - @tag :slow - test "deletes data from events and sessions tables" do - site = insert(:site) - deleted_site = insert(:site) - - populate_stats(site, [ - build(:pageview), - build(:pageview, timestamp: ~D[2026-01-01]), - build(:imported_visitors), - build(:imported_sources), - build(:imported_pages), - build(:imported_entry_pages), - build(:imported_exit_pages), - build(:imported_locations), - build(:imported_devices), - build(:imported_browsers), - build(:imported_operating_systems), - build(:imported_custom_events) - ]) - - populate_stats(deleted_site, [ - build(:pageview), - build(:pageview, timestamp: ~D[2026-01-01]), - build(:pageview, timestamp: ~D[2026-02-01]), - build(:pageview, timestamp: ~D[2026-03-01]), - build(:imported_visitors), - build(:imported_sources), - build(:imported_pages), - build(:imported_entry_pages), - build(:imported_exit_pages), - build(:imported_locations), - build(:imported_devices), - build(:imported_browsers), - build(:imported_operating_systems), - build(:imported_custom_events) - ]) - - Repo.delete!(deleted_site) - - assert Enum.member?( - ClickhouseCleanSites.get_deleted_sites_with_clickhouse_data(), - deleted_site.id - ) - - assert not Enum.member?( - ClickhouseCleanSites.get_deleted_sites_with_clickhouse_data(), - site.id - ) - - ClickhouseCleanSites.perform(nil) - - assert_count(deleted_site, "events_v2", 0) - assert_count(deleted_site, "sessions_v2", 0) - assert_count(deleted_site, "imported_visitors", 0) - assert_count(deleted_site, "imported_sources", 0) - assert_count(deleted_site, "imported_pages", 0) - assert_count(deleted_site, "imported_entry_pages", 0) - assert_count(deleted_site, "imported_exit_pages", 0) - assert_count(deleted_site, "imported_locations", 0) - assert_count(deleted_site, "imported_devices", 0) - assert_count(deleted_site, "imported_browsers", 0) - assert_count(deleted_site, "imported_operating_systems", 0) - assert_count(deleted_site, "imported_custom_events", 0) - assert_count(site, "events_v2", 2) - assert_count(site, "sessions_v2", 2) - assert_count(site, "imported_visitors", 1) - assert_count(site, "imported_sources", 1) - assert_count(site, "imported_pages", 1) - assert_count(site, "imported_entry_pages", 1) - assert_count(site, "imported_exit_pages", 1) - assert_count(site, "imported_locations", 1) - assert_count(site, "imported_devices", 1) - assert_count(site, "imported_browsers", 1) - assert_count(site, "imported_operating_systems", 1) - assert_count(site, "imported_custom_events", 1) - - assert not Enum.member?( - ClickhouseCleanSites.get_deleted_sites_with_clickhouse_data(), - deleted_site.id - ) - end - - @tag :slow - test "cleans a deleted site that has only imported data and no native events" do - kept_site = insert(:site) - imported_only = insert(:site) - - # A live site with native events - must be left untouched. - populate_stats(kept_site, [build(:pageview)]) - - populate_stats(imported_only, [ - build(:imported_visitors), - build(:imported_sources), - build(:imported_custom_events) - ]) - - Repo.delete!(imported_only) - - assert Enum.member?( - ClickhouseCleanSites.get_deleted_sites_with_clickhouse_data(), - imported_only.id - ) - - ClickhouseCleanSites.perform(nil) - - assert_count(imported_only, "imported_visitors", 0) - assert_count(imported_only, "imported_sources", 0) - assert_count(imported_only, "imported_custom_events", 0) - assert_count(kept_site, "events_v2", 1) - end - - def assert_count(site, table, expected_count) do - q = from(e in table, select: %{count: fragment("count()")}, where: e.site_id == ^site.id) - assert await_clickhouse_count(q, expected_count) - end -end