Skip to content
Open
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
7 changes: 7 additions & 0 deletions rel/overlay/etc/default.ini
Original file line number Diff line number Diff line change
Expand Up @@ -724,6 +724,13 @@ partitioned||* = true
; *.example.com:443:[2001:db8::1]:443
;connect_to =

; Compress outbound replication request bodies (_bulk_docs, _revs_diff, _bulk_get).
; Accepted values: none (default, disabled), gzip.
; Enable gzip only when the target supports Content-Encoding: gzip on inbound requests.
;request_compression = none
;compress_min_size = 1024


; Some socket options that might boost performance in some scenarios:
; {nodelay, boolean()}
; {sndbuf, integer()}
Expand Down
7 changes: 6 additions & 1 deletion src/couch_replicator/include/couch_replicator_api_wrap.hrl
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,10 @@

-define(COUCH_REPLICATOR_USER_AGENT, "CouchDB-Replicator/" ++ couch_server:get_version()).

-define(COMPRESS_NONE, "none").
-define(COMPRESS_GZIP, "gzip").
-define(COMPRESS_MIN_SIZE, 1024).

-record(httpdb, {
url,
auth_props = [],
Expand All @@ -28,5 +32,6 @@
http_connections,
first_error_timestamp = nil,
proxy_url,
auth_context = nil
auth_context = nil,
request_compression = ?COMPRESS_NONE
}).
5 changes: 5 additions & 0 deletions src/couch_replicator/priv/stats_descriptions.cfg
Original file line number Diff line number Diff line change
Expand Up @@ -146,3 +146,8 @@
{type, counter},
{desc, <<"number of times DNS overrides were applied to replication requests">>}
]}.

{[couch_replicator, requests_compressed, gzip], [
{type, counter},
{desc, <<"number of HTTP requests compressed with gzip by the replicator">>}
]}.
49 changes: 43 additions & 6 deletions src/couch_replicator/src/couch_replicator_api_wrap.erl
Original file line number Diff line number Diff line change
Expand Up @@ -171,13 +171,15 @@ ensure_full_commit(#httpdb{} = Db) ->

get_missing_revs(#httpdb{} = Db, IdRevs) ->
JsonBody = {[{Id, couch_doc:revs_to_strs(Revs)} || {Id, Revs} <- IdRevs]},
RawBody = ?JSON_ENCODE(JsonBody),
{Body, ExtraHeaders} = maybe_compress_request(Db, RawBody),
send_req(
Db,
[
{method, post},
{path, "_revs_diff"},
{body, ?JSON_ENCODE(JsonBody)},
{headers, [{"Content-Type", "application/json"}]}
{body, Body},
{headers, [{"Content-Type", "application/json"} | ExtraHeaders]}
],
fun
(200, _, {Props}) ->
Expand Down Expand Up @@ -211,14 +213,17 @@ bulk_get(#httpdb{} = Db, #{} = IdRevs, Options) ->
% that at some point in the future we could make that the default, instead
% of having to send query parameters with a POST request as we do today
Body = options_to_json_map(Options, #{<<"docs">> => ReqDocsMaps}),
RawBody = ?JSON_ENCODE(Body),
{ReqBody, ExtraHeaders} = maybe_compress_request(Db, RawBody),
Req = [
{method, post},
{path, "_bulk_get"},
{qs, options_to_query_args(Options, [])},
{body, ?JSON_ENCODE(Body)},
{body, ReqBody},
{headers, [
{"Content-Type", "application/json"},
{"Accept", "application/json"}
| ExtraHeaders
]}
],
try
Expand Down Expand Up @@ -500,17 +505,24 @@ update_docs(#httpdb{} = HttpDb, DocList, Options, UpdateType) ->
([Doc | RestDocs]) ->
{ok, [Doc, ","], RestDocs}
end,
Headers = [
{"Content-Length", Len},
Headers0 = [
{"Content-Type", "application/json"},
{"X-Couch-Full-Commit", FullCommit}
],
{Body, Headers} =
case should_compress_request(HttpDb, Len) of
true ->
FullBody = iolist_to_binary([Prefix, lists:join(",", Docs), Suffix]),

@nickva nickva Aug 6, 2026

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.

Good catch by @rnewson that zlib:gzip takes iolists that means we don't have have to create a binary here. The FullBody could just be [Prefix, lists:join(",", Docs), Suffix] just check that we don't do byte_size or is_binary calls on the result

gzip_request_body(FullBody, Headers0);
false ->
{{BodyFun, [prefix | Docs]}, [{"Content-Length", Len} | Headers0]}
end,
send_req(
HttpDb,
[
{method, post},
{path, "_bulk_docs"},
{body, {BodyFun, [prefix | Docs]}},
{body, Body},
{headers, Headers}
],
fun
Expand Down Expand Up @@ -1052,6 +1064,31 @@ header_value(Key, Headers, Default) ->
_ ->
Default
end.

%% Returns true if compression is enabled and body meets the minimum size threshold.
should_compress_request(#httpdb{request_compression = ?COMPRESS_GZIP}, BodySize) ->
MinSize = config:get_integer("replicator", "compress_min_size", ?COMPRESS_MIN_SIZE),
BodySize >= MinSize;
should_compress_request(#httpdb{}, _BodySize) ->
false.

%% Compress Body with gzip, prepend Content-Length and Content-Encoding headers.
%% Returns {CompressedBody, Headers}.
gzip_request_body(Body, Headers) ->
Compressed = zlib:gzip(Body),
Len = byte_size(Compressed),
couch_stats:increment_counter([couch_replicator, requests_compressed, gzip]),
{Compressed, [{"Content-Length", Len}, {"Content-Encoding", "gzip"} | Headers]}.

%% Compress Body if compression is enabled and body meets minimum size.
%% Returns {Body, ExtraHeaders} where ExtraHeaders may contain Content-Encoding.
maybe_compress_request(#httpdb{} = HttpDb, Body) ->
case should_compress_request(HttpDb, iolist_size(Body)) of
true ->
gzip_request_body(Body, []);
false ->
{Body, []}
end.

% Normalize an #httpdb{} or #db{} record such that it can be used for
% comparisons. This means remove things like pids and also sort options / props.
Expand Down
11 changes: 10 additions & 1 deletion src/couch_replicator/src/couch_replicator_parse.erl
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,7 @@ default_options() ->
{checkpoint_interval, cfg_int("checkpoint_interval", 30000)},
{use_checkpoints, cfg_boolean("use_checkpoints", true)},
{use_bulk_get, cfg_boolean("use_bulk_get", true)},
{request_compression, cfg_str("request_compression", "none")},
{ibrowse_options, cfg_ibrowse_opts()},
{socket_options, cfg_sock_opts()}
].
Expand Down Expand Up @@ -216,7 +217,8 @@ parse_rep_db({Props}, Proxy, Options) ->
timeout = get_value(connection_timeout, Options),
http_connections = get_value(http_connections, Options),
retries = get_value(retries, Options),
proxy_url = ProxyURL
proxy_url = ProxyURL,
request_compression = get_value(request_compression, Options, "none")
},
couch_replicator_utils:normalize_basic_auth(HttpDb);
parse_rep_db(<<"http://", _/binary>> = Url, Proxy, Options) ->
Expand Down Expand Up @@ -284,6 +286,9 @@ cfg_int(Var, Default) ->
cfg_boolean(Var, Default) ->
config:get_boolean("replicator", Var, Default).

cfg_str(Var, Default) ->
config:get("replicator", Var, Default).

cfg_atoms(Cfg, Default) ->
case cfg(Cfg) of
undefined ->
Expand Down Expand Up @@ -388,6 +393,8 @@ convert_options([{<<"since_seq">>, V} | R]) ->
[{since_seq, V} | convert_options(R)];
convert_options([{<<"use_checkpoints">>, V} | R]) ->
[{use_checkpoints, V} | convert_options(R)];
convert_options([{<<"request_compression">>, V} | R]) when is_binary(V) ->
[{request_compression, binary_to_list(V)} | convert_options(R)];
convert_options([{<<"use_bulk_get">>, V} | _R]) when not is_boolean(V) ->
throw({bad_request, <<"parameter `use_bulk_get` must be a boolean">>});
convert_options([{<<"use_bulk_get">>, V} | R]) ->
Expand Down Expand Up @@ -774,6 +781,7 @@ t_parse_sock_opts(_) ->
{connection_timeout, 30000},
{http_connections, 20},
{ibrowse_options, []},
{request_compression, "none"},
{retries, 5},
{socket_options, [
{priority, 3},
Expand Down Expand Up @@ -819,6 +827,7 @@ t_parse_ibrowse_opts(_) ->
{ibrowse_options, [
{prefer_ipv6, true}
]},
{request_compression, "none"},
{retries, 5},
{socket_options, [
{keepalive, true},
Expand Down
126 changes: 126 additions & 0 deletions src/couch_replicator/test/eunit/couch_replicator_compression_tests.erl
Original file line number Diff line number Diff line change
@@ -0,0 +1,126 @@
% Licensed under the Apache License, Version 2.0 (the "License"); you may not
% use this file except in compliance with the License. You may obtain a copy of
% the License at
%
% http://www.apache.org/licenses/LICENSE-2.0
%
% Unless required by applicable law or agreed to in writing, software
% distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
% WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
% License for the specific language governing permissions and limitations under
% the License.

-module(couch_replicator_compression_tests).

-include_lib("couch/include/couch_eunit.hrl").
-include_lib("couch/include/couch_db.hrl").

-define(DOCS_COUNT, 10).
Comment thread
nickva marked this conversation as resolved.
-define(LARGE_DOCS_COUNT, 500).
-define(TIMEOUT_EUNIT, 60).

compression_test_() ->
{
"Replication compression tests",
{
foreach,
fun setup/0,
fun teardown/1,
[
?TDEF_FE(should_not_compress_by_default, ?TIMEOUT_EUNIT),
?TDEF_FE(should_compress_when_enabled, ?TIMEOUT_EUNIT),
?TDEF_FE(should_compress_large_batch, ?TIMEOUT_EUNIT),
?TDEF_FE(should_compress_per_job, ?TIMEOUT_EUNIT),
?TDEF_FE(job_compression_overrides_global_disabled, ?TIMEOUT_EUNIT)
]
}
}.

setup() ->
Ctx = couch_replicator_test_helper:test_setup(),
config:set("replicator", "request_compression", "none", false),
config:set("replicator", "compress_min_size", "1024", false),
Ctx.

teardown(Ctx) ->
config:delete("replicator", "request_compression", false),
config:delete("replicator", "compress_min_size", false),
couch_replicator_test_helper:test_teardown(Ctx).

should_not_compress_by_default({_Ctx, {Source, Target}}) ->
Before = couch_stats:sample([couch_replicator, requests_compressed, gzip]),
populate_db(Source, ?DOCS_COUNT),
replicate(Source, Target),
couch_replicator_test_helper:cluster_compare_dbs(Source, Target),
After = couch_stats:sample([couch_replicator, requests_compressed, gzip]),
?assertEqual(Before, After).

should_compress_when_enabled({_Ctx, {Source, Target}}) ->
config:set("replicator", "request_compression", "gzip", false),
config:set("replicator", "compress_min_size", "10", false),
Before = couch_stats:sample([couch_replicator, requests_compressed, gzip]),
populate_db(Source, ?DOCS_COUNT),
replicate(Source, Target),
couch_replicator_test_helper:cluster_compare_dbs(Source, Target),
After = couch_stats:sample([couch_replicator, requests_compressed, gzip]),
?assert(After > Before).

should_compress_large_batch({_Ctx, {Source, Target}}) ->
config:set("replicator", "request_compression", "gzip", false),
Before = couch_stats:sample([couch_replicator, requests_compressed, gzip]),
populate_db(Source, ?LARGE_DOCS_COUNT),
replicate(Source, Target),
couch_replicator_test_helper:cluster_compare_dbs(Source, Target),
After = couch_stats:sample([couch_replicator, requests_compressed, gzip]),
?assert(After > Before).

should_compress_per_job({_Ctx, {Source, Target}}) ->
% global config is none (default), but job sets gzip
config:set("replicator", "compress_min_size", "10", false),
Before = couch_stats:sample([couch_replicator, requests_compressed, gzip]),
populate_db(Source, ?DOCS_COUNT),
replicate_with_options(Source, Target, [{<<"request_compression">>, <<"gzip">>}]),
couch_replicator_test_helper:cluster_compare_dbs(Source, Target),
After = couch_stats:sample([couch_replicator, requests_compressed, gzip]),
?assert(After > Before).

job_compression_overrides_global_disabled({_Ctx, {Source, Target}}) ->
% global config is gzip, but job disables it
config:set("replicator", "request_compression", "gzip", false),
config:set("replicator", "compress_min_size", "10", false),
Before = couch_stats:sample([couch_replicator, requests_compressed, gzip]),
populate_db(Source, ?DOCS_COUNT),
replicate_with_options(Source, Target, [{<<"request_compression">>, <<"none">>}]),
couch_replicator_test_helper:cluster_compare_dbs(Source, Target),
After = couch_stats:sample([couch_replicator, requests_compressed, gzip]),
?assertEqual(Before, After).

populate_db(DbName, Count) ->
Docs = lists:map(
fun(I) ->
Id = iolist_to_binary(io_lib:format("doc~p", [I])),
Data = list_to_binary(lists:duplicate(100, $x)),
{[
{<<"_id">>, Id},
{<<"value">>, I},
{<<"data">>, Data}
]}
end,
lists:seq(1, Count)
),
{ok, _} = fabric:update_docs(DbName, Docs, [?ADMIN_CTX]),
ok.

replicate(Source, Target) ->
replicate_with_options(Source, Target, []).

replicate_with_options(Source, Target, ExtraOptions) ->
SourceUrl = couch_replicator_test_helper:cluster_db_url(Source),
TargetUrl = couch_replicator_test_helper:cluster_db_url(Target),
RepObject = {[
{<<"source">>, SourceUrl},
{<<"target">>, TargetUrl},
{<<"continuous">>, false}
| ExtraOptions
]},
{ok, _} = couch_replicator_test_helper:replicate(RepObject).