Skip to content
Draft
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
43 changes: 43 additions & 0 deletions bench/log_event_cleaner.exs
Original file line number Diff line number Diff line change
@@ -0,0 +1,43 @@
# Run from the repository root with:
# MIX_ENV=dev ../bin/x mix run --no-start bench/log_event_cleaner.exs
#
# Compare this command unchanged on the benchmark-only baseline commit and its
# copy-on-write cleaner child.

alias Logflare.Logs.IngestTransformers

benchmark_time = String.to_integer(System.get_env("LF_BENCH_TIME") || "5")
warmup_time = String.to_integer(System.get_env("LF_BENCH_WARMUP") || "2")
memory_time = String.to_integer(System.get_env("LF_BENCH_MEMORY_TIME") || "3")

clean = Map.new(1..64, &{"field_#{&1}", &1})
sparse_dirty = clean |> Map.put("drop", nil) |> Map.put("bad-key", true)
dense_empty = Map.new(1..64, &{"field_#{&1}", nil})
dense_unsafe = Map.new(1..64, &{"field-#{&1}", &1})
half_empty = Map.new(1..64, &{"field_#{&1}", if(rem(&1, 2) == 0, do: nil, else: &1)})

list_clean = %{"items" => Enum.map(1..64, &%{"id" => &1, "value" => "ok"})}
list_sparse_dirty = put_in(list_clean, ["items", Access.at(32), "drop"], nil)
list_half_empty = %{"items" => Enum.flat_map(1..32, &[nil, %{"id" => &1}])}

scenarios = %{
"map: clean 64" => clean,
"map: sparse dirty 66" => sparse_dirty,
"map: half empty 64" => half_empty,
"map: all empty 64" => dense_empty,
"map: all unsafe 64" => dense_unsafe,
"list: clean 64 maps" => list_clean,
"list: one nested change" => list_sparse_dirty,
"list: half empty" => list_half_empty
}

Benchee.run(
Map.new(scenarios, fn {name, payload} ->
{name, fn -> IngestTransformers.transform(payload, :clean_to_bigquery_column_spec) end}
end),
time: benchmark_time,
warmup: warmup_time,
memory_time: memory_time,
parallel: 1,
print: [benchmarking: true, configuration: true, fast_warning: false]
)
108 changes: 89 additions & 19 deletions lib/logflare/logs/ingest_transformers.ex
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,8 @@ defmodule Logflare.Logs.IngestTransformers do
alias Logflare.Logs.Ingest.MetadataCleaner

@alphanumeric_regex ~r/\W/
# Above this size, scanning an all-empty map allocates less than deleting every key.
@all_empty_map_threshold 16
@max_field_length 128

@typep direct_transform :: :clean_to_bigquery_column_spec | :to_bigquery_column_spec
Expand Down Expand Up @@ -38,46 +40,94 @@ defmodule Logflare.Logs.IngestTransformers do
Enum.reduce(rules, log_params, &do_transform(&2, &1))
end

defp clean_and_to_bigquery_column_spec(%_{} = data), do: rebuild_clean_map(data)

defp clean_and_to_bigquery_column_spec(data) when is_map(data) do
if map_size(data) >= @all_empty_map_threshold and all_empty_map?(data) do
%{}
else
copy_on_write_clean_map(data)
end
end

defp clean_and_to_bigquery_column_spec(data) when is_list(data) do
data
|> Enum.reduce([], fn
value, acc when is_nil_or_empty(value) ->
acc

value, acc when is_map(value) or is_list(value) ->
cleaned = clean_and_to_bigquery_column_spec(value)
if is_nil_or_empty(cleaned), do: acc, else: [cleaned | acc]

value, acc ->
[value | acc]
end)
|> Enum.reverse()
end

defp copy_on_write_clean_map(data) do
:maps.fold(
fn
_k, v, acc when is_nil_or_empty(v) ->
key, value, acc when is_nil_or_empty(value) ->
Map.delete(acc, key)

key, value, acc when is_map(value) or is_list(value) ->
cleaned = clean_and_to_bigquery_column_spec(value)

if is_nil_or_empty(cleaned) do
Map.delete(acc, key)
else
update_bigquery_column(acc, key, value, cleaned)
end

key, value, acc ->
update_bigquery_column(acc, key, value, value)
end,
data,
data
)
end

defp rebuild_clean_map(data) do
:maps.fold(
fn
_key, value, acc when is_nil_or_empty(value) ->
acc

k, v, acc when is_map(v) or is_list(v) ->
cleaned = clean_and_to_bigquery_column_spec(v)
key, value, acc when is_map(value) or is_list(value) ->
cleaned = clean_and_to_bigquery_column_spec(value)

if is_nil_or_empty(cleaned) do
acc
else
put_bigquery_column(acc, k, cleaned)
put_bigquery_column(acc, key, cleaned)
end

k, v, acc ->
put_bigquery_column(acc, k, v)
key, value, acc ->
put_bigquery_column(acc, key, value)
end,
%{},
data
)
end

defp clean_and_to_bigquery_column_spec(data) when is_list(data) do
data
|> Enum.reduce([], fn
x, acc when is_nil_or_empty(x) ->
acc
defp all_empty_map?(data), do: all_empty_map_iterator?(:maps.iterator(data))

x, acc when is_map(x) or is_list(x) ->
cleaned = clean_and_to_bigquery_column_spec(x)
if is_nil_or_empty(cleaned), do: acc, else: [cleaned | acc]
defp all_empty_map_iterator?(iterator) do
case :maps.next(iterator) do
{_key, value, next_iterator} when is_nil_or_empty(value) ->
all_empty_map_iterator?(next_iterator)

x, acc ->
[x | acc]
end)
|> Enum.reverse()
:none ->
true

_entry ->
false
end
end

@compile {:inline, put_bigquery_column: 3}
@compile {:inline, put_bigquery_column: 3, update_bigquery_column: 4}
@spec put_bigquery_column(map(), term(), term()) :: map()
defp put_bigquery_column(acc, key, value) do
result = Map.put(acc, to_bigquery_column_spec(key), value)
Expand All @@ -87,6 +137,26 @@ defmodule Logflare.Logs.IngestTransformers do
result
end

@spec update_bigquery_column(map(), term(), term(), term()) :: map()
defp update_bigquery_column(acc, key, original_value, cleaned_value) do
case to_bigquery_column_spec(key) do
^key when original_value == cleaned_value ->
acc

^key ->
Map.put(acc, key, cleaned_value)

normalized_key ->
acc = Map.delete(acc, key)

if Map.has_key?(acc, normalized_key) do
throw(:normalized_bigquery_column_collision)
end

Map.put(acc, normalized_key, cleaned_value)
end
end

# Rewrites a map key into a valid BigQuery standard column name in a single
# pass. Standard names allow only [A-Za-z0-9_], cannot start with a digit,
# cannot use a reserved prefix, and must be valid UTF-8. Classification is per
Expand Down
51 changes: 51 additions & 0 deletions test/logflare/logs/ingest/ingest_transformer_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -418,6 +418,57 @@ defmodule Logflare.Logs.IngestTransformerTest do
end

describe ":clean_to_bigquery_column_spec fused pipeline" do
test "preserves clean nested payload values" do
input = %{
"safe" => %{"nested" => 1, "values" => [true, 2, "three"]},
"items" => [%{"id" => 1}, %{"id" => 2}],
:non_binary_key => "kept"
}

assert transform(input, :clean_to_bigquery_column_spec) == input
end

test "cleans sparse map changes without altering unaffected values" do
input = %{
"safe" => %{"kept" => 1, "removed" => nil},
"bad-key" => %{"nested" => true},
"items" => [nil, %{"kept" => 2, "removed" => ""}, false]
}

assert transform(input, :clean_to_bigquery_column_spec) == %{
"safe" => %{"kept" => 1},
"_bad_key" => %{"nested" => true},
"items" => [%{"kept" => 2}, false]
}
end

test "removes large maps containing only empty values" do
input = Map.new(1..16, &{"empty_#{&1}", nil})
assert transform(input, :clean_to_bigquery_column_spec) == %{}
end

test "normalizes every key in large unsafe-key maps" do
input = Map.new(1..16, &{"field-#{&1}", &1})

expected =
input
|> MetadataCleaner.deep_reject_nil_and_empty()
|> transform(:to_bigquery_column_spec)

assert transform(input, :clean_to_bigquery_column_spec) == expected
end

test "retains sequential struct-cleaning behavior" do
input = %{"date" => ~D[2026-09-13]}

expected =
input
|> MetadataCleaner.deep_reject_nil_and_empty()
|> transform(:to_bigquery_column_spec)

assert transform(input, :clean_to_bigquery_column_spec) == expected
end

property "matches sequential cleaning and key normalization" do
check all input <- log_params_generator() do
assert transform(input, :clean_to_bigquery_column_spec) ==
Expand Down
Loading