Skip to content
Merged
2 changes: 1 addition & 1 deletion VERSION
Original file line number Diff line number Diff line change
@@ -1 +1 @@
1.50.6
1.50.7
4 changes: 1 addition & 3 deletions config/config.exs
Original file line number Diff line number Diff line change
Expand Up @@ -34,9 +34,7 @@ config :logflare, :bigquery_backend_adaptor, managed_service_account_pool_size:

config :logflare, :bigquery_pipeline, max_retries: 0

config :logflare, :clickhouse_backend_adaptor,
engine: "MergeTree",
pool_size: 3
config :logflare, :clickhouse_backend_adaptor, engine: "MergeTree"

config :logflare, Logflare.Sources.Source.BigQuery.Schema, updates_per_minute: 6

Expand Down
55 changes: 36 additions & 19 deletions lib/logflare/backends/adaptor/clickhouse_adaptor.ex
Original file line number Diff line number Diff line change
Expand Up @@ -234,7 +234,8 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor do
password: :string,
database: :string,
port: :integer,
pool_size: :integer,
read_pool_size: :integer,
labeled_read_pool_size: :integer,
# read_only_url is depreciated and will be removed in the release after PR#3693 lands
read_only_url: :string,
read_only_urls: {:map, :string},
Expand All @@ -250,7 +251,8 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor do
:password,
:database,
:port,
:pool_size,
:read_pool_size,
:labeled_read_pool_size,
:read_only_url,
:read_only_urls,
:default_read_cluster,
Expand Down Expand Up @@ -303,7 +305,11 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor do
|> validate_read_only_urls()
|> validate_default_read_cluster()
|> validate_user_pass()
|> validate_number(:pool_size,
|> validate_number(:read_pool_size,
greater_than_or_equal_to: 1,
less_than_or_equal_to: @max_read_pool_size
)
|> validate_number(:labeled_read_pool_size,
greater_than_or_equal_to: 1,
less_than_or_equal_to: @max_read_pool_size
)
Expand Down Expand Up @@ -403,8 +409,8 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor do
Logger.warning(
"ClickHouse read cluster GRANT check failed for #{target}. Required: `SELECT`",
backend_id: backend.id,
read_cluster: label,
read_cluster_url: url
clickhouse_read_cluster: label,
clickhouse_read_cluster_url: url
)

{:error, :read_permissions_missing}
Expand All @@ -413,8 +419,8 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor do
Logger.warning(
"ClickHouse read cluster connection/GRANT check failed for #{target}. Unexpected error #{inspect(error_result)}",
backend_id: backend.id,
read_cluster: label,
read_cluster_url: url
clickhouse_read_cluster: label,
clickhouse_read_cluster_url: url
)

{:error, :grant_check_unknown_failure}
Expand Down Expand Up @@ -574,7 +580,7 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor do
timeout = if Application.get_env(:logflare, :env) == :test, do: 1_000, else: 60_000

backend_id = backend.id
log_fun = fn entry -> log_slow_checkout(entry, backend_id) end
log_fun = fn entry -> log_slow_checkout(entry, backend_id, label) end

ch_opts = [decode: false, timeout: timeout, log: log_fun]

Expand All @@ -592,7 +598,7 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor do
user_id: backend.user_id,
backend_id: backend.id,
backend_token: backend.token,
read_cluster: label,
clickhouse_read_cluster: label,
host: ConnectionManager.read_host(backend, label)
)}
end
Expand All @@ -606,8 +612,8 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor do
"ClickHouse read cluster not configured, falling back to resolved read cluster",
user_id: backend.user_id,
backend_id: backend.id,
requested_read_cluster: requested,
resolved_read_cluster: label
clickhouse_requested_read_cluster: requested,
clickhouse_resolved_read_cluster: label
)

:ok
Expand All @@ -630,8 +636,8 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor do
"ClickHouse read cluster unhealthy, falling back to default read cluster",
user_id: backend.user_id,
backend_id: backend.id,
read_cluster: label,
default_read_cluster: default
clickhouse_read_cluster: label,
clickhouse_default_read_cluster: default
)

do_ch_query_on_label(backend, statement, params, default)
Expand All @@ -640,22 +646,29 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor do
end
end

@spec log_slow_checkout(DBConnection.LogEntry.t(), pos_integer()) :: :ok
defp log_slow_checkout(%DBConnection.LogEntry{pool_time: pool_time}, backend_id)
@spec log_slow_checkout(DBConnection.LogEntry.t(), pos_integer(), String.t() | nil) :: :ok
defp log_slow_checkout(%DBConnection.LogEntry{pool_time: pool_time}, backend_id, label)
when is_integer(pool_time) do
pool_ms = System.convert_time_unit(pool_time, :native, :millisecond)

:telemetry.execute(
[:logflare, :clickhouse, :read_pool, :checkout],
%{pool_time_ms: pool_ms},
%{backend_id: backend_id, read_cluster: label}
)

if pool_ms >= slow_pool_checkout_ms() do
Logger.warning(
"ClickHouse slow connection checkout: waited #{pool_ms}ms for a pool connection",
backend_id: backend_id
backend_id: backend_id,
clickhouse_read_cluster: label
)
end

:ok
end

defp log_slow_checkout(_entry, _backend_id), do: :ok
defp log_slow_checkout(_entry, _backend_id, _label), do: :ok

@spec slow_pool_checkout_ms() :: non_neg_integer()
defp slow_pool_checkout_ms do
Expand All @@ -670,6 +683,10 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor do
|> query_error(error)
end

defp to_query_error(%DBConnection.ConnectionError{reason: :queue_timeout} = error) do
query_error(:pool_exhausted, error)
end

defp to_query_error(%DBConnection.ConnectionError{} = error) do
query_error(:connection_error, error)
end
Expand Down Expand Up @@ -1240,7 +1257,7 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor do
Logger.info(
"Started query ConnectionManager for ClickHouse backend",
backend_id: backend.id,
read_cluster: label
clickhouse_read_cluster: label
)

:ok
Expand All @@ -1252,7 +1269,7 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor do
Logger.warning(
"Failed to start query ConnectionManager for backend",
backend_id: backend_id,
read_cluster: label,
clickhouse_read_cluster: label,
reason: reason
)

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,9 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor.ConnectionManager do
@ch_queue_target :timer.seconds(5)
@recycle_interval :timer.minutes(10)
@recycle_spread :timer.seconds(60)
@ch_idle_interval :timer.seconds(3)
@default_read_pool_size 50
@default_labeled_read_pool_size 32

typedstruct do
field :backend_id, pos_integer(), enforce: true
Expand Down Expand Up @@ -211,6 +214,24 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor.ConnectionManager do

def read_host(_backend, _label), do: nil

@doc """
Resolves the read connection pool size for a labeled read cluster.

The unlabeled pool and the pool for the cluster named by `default_read_cluster`
are sized by `read_pool_size`. Both act as the catch-all for callers that do not
resolve to a read cluster of their own, and the default cluster additionally
absorbs failover traffic, so they need more capacity than a single caller's
cluster. Every other labeled pool is sized by `labeled_read_pool_size`.

A backend-wide value wins over the application default.
"""
@spec read_pool_size(map(), String.t() | nil) :: pos_integer()
def read_pool_size(config, label) when is_map(config) do
config
|> pool_size_key(label)
|> resolve_pool_size(config)
end

@impl true
def init({backend_id, label}) do
resolve_timer_ref = resolve_timer_send_after()
Expand Down Expand Up @@ -462,16 +483,7 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor.ConnectionManager do

defp build_ch_opts(%Backend{} = backend, label) do
config = backend.config

default_pool_size =
Application.fetch_env!(:logflare, :clickhouse_backend_adaptor)[:pool_size]

pool_size =
config
|> Map.get(:pool_size, default_pool_size)
|> div(2)
|> max(default_pool_size)

pool_size = read_pool_size(config, label)
url = read_url(config, label)

with {:ok, {scheme, hostname, url_port}} <- extract_url_components(url) do
Expand All @@ -489,13 +501,41 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor.ConnectionManager do
pool_size: pool_size,
settings: [],
timeout: @ch_query_conn_timeout,
queue_target: @ch_queue_target
queue_target: @ch_queue_target,
idle_interval: @ch_idle_interval,
idle_limit: pool_size
]

{:ok, ch_opts}
end
end

@spec pool_size_key(map(), String.t() | nil) :: :read_pool_size | :labeled_read_pool_size
defp pool_size_key(_config, label) when not is_non_empty_binary(label), do: :read_pool_size

defp pool_size_key(config, label) do
case Map.get(config, :default_read_cluster) do
^label -> :read_pool_size
_ -> :labeled_read_pool_size
end
end

@spec resolve_pool_size(atom(), map()) :: pos_integer()
defp resolve_pool_size(key, config) do
case validate_pool_size(Map.get(config, key)) do
nil -> default_pool_size(key)
size -> size
end
end

@spec default_pool_size(atom()) :: pos_integer()
defp default_pool_size(:read_pool_size), do: @default_read_pool_size
defp default_pool_size(:labeled_read_pool_size), do: @default_labeled_read_pool_size

@spec validate_pool_size(term()) :: pos_integer() | nil
defp validate_pool_size(size) when is_pos_integer(size), do: size
defp validate_pool_size(_size), do: nil

@spec read_url(map(), String.t() | nil) :: String.t() | nil
defp read_url(config, label) do
urls = Map.get(config, :read_only_urls) || %{}
Expand Down
3 changes: 2 additions & 1 deletion lib/logflare/backends/query_error.ex
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,8 @@ defmodule Logflare.Backends.QueryError do
@enforce_keys [:kind, :raw_error, :backend]
defstruct [:kind, :raw_error, :backend, :description]

@type kind :: :invalid_query | :connection_error | :backend_error | :timeout
@type kind ::
:invalid_query | :connection_error | :pool_exhausted | :backend_error | :timeout
@type t :: %__MODULE__{
kind: kind(),
raw_error: term(),
Expand Down
7 changes: 5 additions & 2 deletions lib/logflare_web/live/backends/components/backend_form.heex
Original file line number Diff line number Diff line change
Expand Up @@ -258,8 +258,11 @@
</small>
</div>
<div class="form-group">
{label(f_config, :pool_size, "HTTP Pool Size")}
{text_input(f_config, :pool_size, class: "form-control", value: input_value(f_config, "pool_size") || 20)}
{label(f_config, :read_pool_size, "Read Pool Size")}
{text_input(f_config, :read_pool_size, class: "form-control", value: input_value(f_config, "read_pool_size") || 50)}
<small class="form-text text-muted">
Connections per node for the primary read pool.
</small>
</div>
<div class="form-group">
{label(f_config, :read_only_url, "Read-Only URL (Optional)")}
Expand Down
33 changes: 32 additions & 1 deletion lib/logflare_web/live/backends/read_cluster_urls_component.ex
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,17 @@ defmodule LogflareWeb.Backends.ReadClusterUrlsComponent do
def render(assigns) do
~H"""
<div>
<div class="form-group">
{label(@form, :labeled_read_pool_size, "Read Cluster Pool Size")}
{text_input(@form, :labeled_read_pool_size,
class: "form-control",
value: input_value(@form, "labeled_read_pool_size") || 32
)}
<small class="form-text text-muted">
Connections per node for each read cluster below. The default read cluster uses the primary read pool size instead, since it also absorbs unrecognized callers and failover traffic.
</small>
</div>

<div class="form-group">
<label>Read-Only Cluster URLs (Optional)</label>
<small class="form-text text-muted">
Expand Down Expand Up @@ -111,18 +122,22 @@ defmodule LogflareWeb.Backends.ReadClusterUrlsComponent do
</div>

<div class="form-group">
{label(@form, :default_read_cluster, "Default Read Cluster (Optional)")}
{label(@form, :default_read_cluster, default_read_cluster_label(cluster_configured?(@rows)))}
{text_input(@form, :default_read_cluster,
value: @default_read_cluster,
class: "form-control",
placeholder: "caller label",
required: cluster_configured?(@rows),
Comment thread
amokan marked this conversation as resolved.
phx_change: "sync",
phx_target: @myself,
phx_debounce: "blur"
)}
<small class="form-text text-muted">
The caller label whose cluster absorbs unrecognized or absent callers. Must match a label above.
</small>
<small :if={default_read_cluster_missing?(@rows, @default_read_cluster)} class="form-text tw-text-red-500">
Required once a read cluster is configured. Without it, callers that send no label read from the ingest cluster instead.
</small>
</div>
</div>
"""
Expand Down Expand Up @@ -176,6 +191,22 @@ defmodule LogflareWeb.Backends.ReadClusterUrlsComponent do

defp initial_rows(_backend), do: [{0, "", ""}]

@spec cluster_configured?([row()]) :: boolean()
defp cluster_configured?(rows) do
Enum.any?(rows, fn {_ref, label, url} ->
is_non_empty_binary(label) and is_non_empty_binary(url)
end)
end

@spec default_read_cluster_missing?([row()], String.t() | nil) :: boolean()
defp default_read_cluster_missing?(_rows, default) when is_non_empty_binary(default), do: false

defp default_read_cluster_missing?(rows, _default), do: cluster_configured?(rows)

@spec default_read_cluster_label(boolean()) :: String.t()
defp default_read_cluster_label(true), do: "Default Read Cluster"
defp default_read_cluster_label(false), do: "Default Read Cluster (Optional)"

@spec read_cluster_form_key?(String.t()) :: boolean()
defp read_cluster_form_key?(key) do
String.starts_with?(key, "read_cluster_label_") or
Expand Down
2 changes: 1 addition & 1 deletion lib/logflare_web/live/display_helpers.ex
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ defmodule LogflareWeb.Live.DisplayHelpers do
"""
def sanitize_backend_config(config) when is_map(config) do

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

for future, probably want to let the Adaptor handle this declaration based on their config keys

allowed_keys =
~w(batch_timeout database default_read_cluster hostname host max_event_age_hours max_message_bytes pool_size port project_id read_only_url read_only_urls region s3_bucket schema storage_region table url use_async_inserts_for_small_batches async_insert_cluster_url async_insert_max_rows)a
~w(batch_timeout database default_read_cluster hostname labeled_read_pool_size host max_event_age_hours max_message_bytes port project_id read_only_url read_only_urls read_pool_size region s3_bucket schema storage_region table url use_async_inserts_for_small_batches async_insert_cluster_url async_insert_max_rows)a

config
|> Enum.map(fn {key, value} ->
Expand Down
3 changes: 2 additions & 1 deletion lib/logflare_web/open_api_schemas.ex
Original file line number Diff line number Diff line change
Expand Up @@ -380,7 +380,8 @@ defmodule LogflareWeb.OpenApiSchemas do
port: %Schema{type: :integer},
username: %Schema{type: :string, nullable: true},
password: %Schema{type: :string, nullable: true},
pool_size: %Schema{type: :integer, nullable: true},
read_pool_size: %Schema{type: :integer, nullable: true},
labeled_read_pool_size: %Schema{type: :integer, nullable: true},
read_only_url: %Schema{type: :string, nullable: true},
use_async_inserts_for_small_batches: %Schema{type: :boolean, nullable: true},
async_insert_cluster_url: %Schema{type: :string, nullable: true},
Expand Down
7 changes: 7 additions & 0 deletions lib/telemetry.ex
Original file line number Diff line number Diff line change
Expand Up @@ -256,6 +256,13 @@ defmodule Logflare.Telemetry do
description:
"Sum of events dropped by a backend (timestamp older than its configured max event age)"
),
distribution("logflare.clickhouse.read_pool.checkout",
event_name: [:logflare, :clickhouse, :read_pool, :checkout],
measurement: :pool_time_ms,
unit: :millisecond,
tags: [:backend_id, :read_cluster],
description: "Time spent waiting to check out a ClickHouse read pool connection"
),
sum("logflare.logs.ingest_logs.drop_future",
event_name: [:logflare, :logs, :ingest_logs, :drop_future],
measurement: :count,
Expand Down
Loading
Loading