Packages
electric
1.1.2
1.7.8
1.7.7
1.7.6
1.7.5
1.7.4
1.7.3
1.7.2
1.7.1
1.7.0
1.6.10
1.6.9
1.6.8
1.6.7
1.6.6
1.6.5
1.6.4
1.6.3
1.6.2
1.6.1
1.6.0
1.5.1
1.5.0
1.4.16
1.4.16-beta-1
1.4.15
1.4.14
1.4.13
1.4.12
1.4.11
1.4.10
1.4.8
1.4.7
1.4.6
1.4.5
1.4.4
1.4.3
1.4.2
1.4.1
1.4.0
1.3.4
1.3.3
1.3.2
1.2.4
1.2.3
1.2.2
1.2.1
1.2.0
1.1.14
1.1.13
1.1.12
1.1.11
1.1.10
1.1.9
1.1.8
1.1.7
1.1.6
retired
1.1.5
retired
1.1.4
retired
1.1.3
retired
1.1.2
1.1.1
1.1.0
1.0.24
1.0.23
1.0.22
1.0.21
1.0.20
1.0.19
1.0.18
1.0.17
1.0.15
1.0.13
1.0.12
1.0.11
1.0.10
1.0.9
1.0.5
1.0.4
1.0.3
1.0.2
1.0.1
1.0.0
1.0.0-beta.23
1.0.0-beta.22
1.0.0-beta.20
1.0.0-beta.19
1.0.0-beta.18
1.0.0-beta.17
1.0.0-beta.16
1.0.0-beta.15
1.0.0-beta.14
1.0.0-beta.13
1.0.0-beta.12
1.0.0-beta.11
1.0.0-beta.10
1.0.0-beta.9
1.0.0-beta.8
1.0.0-beta.7
1.0.0-beta.6
1.0.0-beta.5
1.0.0-beta.4
1.0.0-beta.3
1.0.0-beta.2
1.0.0-beta.1
0.9.5
0.9.4
0.9.3
0.9.2
0.9.1
0.9.0
0.8.1
0.8.0
0.7.7
0.7.6
0.7.5
0.7.4
0.7.3
0.7.2
0.7.1
0.7.0
0.6.3
0.6.2
0.6.1
0.5.2
0.4.4
Postgres sync engine. Sync little subsets of your Postgres data into local apps and services.
Current section
Files
Jump to
Current section
Files
lib/electric/shape_cache/file_storage/compaction.ex
defmodule Electric.ShapeCache.FileStorage.Compaction do
alias Electric.LogItems
alias Electric.Utils
alias Electric.ShapeCache.LogChunker
alias Electric.ShapeCache.FileStorage.LogFile
alias Electric.ShapeCache.FileStorage.KeyIndex
alias Electric.ShapeCache.FileStorage.ActionFile
# Compaction and race conditions
#
# `FileStorage` has a pointer to the last compacted offset (and the compacted log file name)
# which is updated atomically once the compaction is complete, so while it's ongoing, the
# pointer is not updated.
#
# While the log is compacted in place, it's actually a merged copy that's being
# compacted, not the original log. Original log is deleted after the compaction
# is complete and the pointer is updated.
#
# Any concurrent reads of the log that's being replaced are also OK: the `File.rename`
# on linux doesn't close the original file descriptor, so the reader will still see
# the original file, and we don't reuse file names. Any readers mid-file of the
# log that's being replaced but that read the chunk will continue from a correct chunk
# of the new file due to offset ordering being preserved. They might observe some updates
# more than once in a compacted form.
@spec compact_in_place({String.t(), String.t(), String.t()}, non_neg_integer(), (any(), any() ->
any())) ::
{String.t(), String.t(), String.t()}
def compact_in_place(
{log_file_path, chunk_index_path, key_index_path},
chunk_size \\ LogChunker.default_chunk_size_threshold(),
merge_fun \\ &LogItems.merge_updates/2
) do
KeyIndex.sort(key_index_path)
ActionFile.create_from_key_index(key_index_path, log_file_path <> ".actions")
{new_log, new_chunk_index, new_key_index} =
LogFile.apply_actions(log_file_path, log_file_path <> ".actions", chunk_size, merge_fun)
File.rm!(log_file_path <> ".actions")
File.rename!(new_log, log_file_path)
File.rename!(new_chunk_index, chunk_index_path)
File.rename!(new_key_index, key_index_path)
{log_file_path, chunk_index_path, key_index_path}
end
def merge_and_compact(
log1,
log2,
merged_log_path,
chunk_size \\ LogChunker.default_chunk_size_threshold()
) do
{log_file_path1, _, key_index_path1} = log1
{log_file_path2, _, key_index_path2} = log2
second_part_start = File.stat!(log_file_path1).size
Utils.concat_files([log_file_path1, log_file_path2], merged_log_path)
KeyIndex.merge_with_offset(
key_index_path1,
key_index_path2,
merged_log_path <> ".key_index",
second_part_start
)
compact_in_place(
{merged_log_path, merged_log_path <> ".chunk_index", merged_log_path <> ".key_index"},
chunk_size
)
end
def rm_log({log, chunk_index, key_index}) do
File.rm!(log)
File.rm!(chunk_index)
File.rm!(key_index)
end
end