diff --git a/rel/overlay/etc/default.ini b/rel/overlay/etc/default.ini index ff148b2714c..53f63a493b0 100644 --- a/rel/overlay/etc/default.ini +++ b/rel/overlay/etc/default.ini @@ -724,6 +724,13 @@ partitioned||* = true ; *.example.com:443:[2001:db8::1]:443 ;connect_to = +; Compress outbound replication request bodies (_bulk_docs, _revs_diff). +; 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()} diff --git a/src/couch_replicator/include/couch_replicator_api_wrap.hrl b/src/couch_replicator/include/couch_replicator_api_wrap.hrl index 6d6ad527cbd..71adb556d60 100644 --- a/src/couch_replicator/include/couch_replicator_api_wrap.hrl +++ b/src/couch_replicator/include/couch_replicator_api_wrap.hrl @@ -28,5 +28,6 @@ http_connections, first_error_timestamp = nil, proxy_url, - auth_context = nil + auth_context = nil, + request_compression = "none" }). diff --git a/src/couch_replicator/priv/stats_descriptions.cfg b/src/couch_replicator/priv/stats_descriptions.cfg index 10821d88516..546b8af38ec 100644 --- a/src/couch_replicator/priv/stats_descriptions.cfg +++ b/src/couch_replicator/priv/stats_descriptions.cfg @@ -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">>} +]}. diff --git a/src/couch_replicator/src/couch_replicator_api_wrap.erl b/src/couch_replicator/src/couch_replicator_api_wrap.erl index 9364757d6cb..5a7a6a6b85e 100644 --- a/src/couch_replicator/src/couch_replicator_api_wrap.erl +++ b/src/couch_replicator/src/couch_replicator_api_wrap.erl @@ -57,6 +57,10 @@ -define(MAX_URL_LEN, 7000). -define(MIN_URL_LEN, 200). +-define(COMPRESS_MIN_SIZE, 1024). +-define(COMPRESS_NONE, "none"). +-define(COMPRESS_GZIP, "gzip"). + db_uri(#httpdb{url = Url}) -> couch_util:url_strip_password(Url). @@ -171,13 +175,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(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}) -> @@ -500,17 +506,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 compress_requests(HttpDb, Len) of + true -> + FullBody = iolist_to_binary([Prefix, lists:join(",", Docs), Suffix]), + 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 @@ -1052,6 +1065,30 @@ header_value(Key, Headers, Default) -> _ -> Default end. + +%% Returns true if compression is enabled for HttpDb and Body is large enough. +compress_requests(#httpdb{request_compression = ?COMPRESS_NONE}, _BodySize) -> + false; +compress_requests(#httpdb{}, BodySize) -> + MinSize = config:get_integer("replicator", "compress_min_size", ?COMPRESS_MIN_SIZE), + BodySize >= MinSize. + +%% Compress Body with gzip, prepend Content-Length and Content-Encoding headers. +%% Returns {CompressedBody, Headers}. +gzip_request_body(Body, Headers) -> + Compressed = zlib:gzip(Body), + couch_stats:increment_counter([couch_replicator, requests_compressed, gzip]), + {Compressed, [{"Content-Length", byte_size(Compressed)}, {"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(#httpdb{} = HttpDb, Body) when is_binary(Body) -> + case compress_requests(HttpDb, byte_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. diff --git a/src/couch_replicator/src/couch_replicator_parse.erl b/src/couch_replicator/src/couch_replicator_parse.erl index 000108acd50..8713a61f416 100644 --- a/src/couch_replicator/src/couch_replicator_parse.erl +++ b/src/couch_replicator/src/couch_replicator_parse.erl @@ -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()} ]. @@ -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) -> @@ -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 -> @@ -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]) -> @@ -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}, @@ -819,6 +827,7 @@ t_parse_ibrowse_opts(_) -> {ibrowse_options, [ {prefer_ipv6, true} ]}, + {request_compression, "none"}, {retries, 5}, {socket_options, [ {keepalive, true}, diff --git a/src/couch_replicator/test/eunit/couch_replicator_compression_tests.erl b/src/couch_replicator/test/eunit/couch_replicator_compression_tests.erl new file mode 100644 index 00000000000..e90e75c3735 --- /dev/null +++ b/src/couch_replicator/test/eunit/couch_replicator_compression_tests.erl @@ -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). +-define(TIMEOUT_EUNIT, 30). + +compression_test_() -> + { + "Replication compression tests", + { + foreach, + fun couch_replicator_test_helper:test_setup/0, + fun couch_replicator_test_helper:test_teardown/1, + [ + ?TDEF_FE(should_not_compress_by_default, ?TIMEOUT_EUNIT), + ?TDEF_FE(should_compress_when_enabled, ?TIMEOUT_EUNIT), + ?TDEF_FE(should_compress_per_job, ?TIMEOUT_EUNIT), + ?TDEF_FE(job_compression_overrides_global_disabled, ?TIMEOUT_EUNIT) + ] + } + }. + +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), + 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), + try + Before = couch_stats:sample([couch_replicator, requests_compressed, gzip]), + populate_db(Source, ?DOCS_COUNT), + replicate(Source, Target), + compare_dbs(Source, Target), + After = couch_stats:sample([couch_replicator, requests_compressed, gzip]), + ?assert(After > Before) + after + config:delete("replicator", "request_compression", false), + config:delete("replicator", "compress_min_size", false) + end. + +should_compress_per_job({_Ctx, {Source, Target}}) -> + % global config is none (default), but job sets gzip + config:set("replicator", "compress_min_size", "10", false), + try + Before = couch_stats:sample([couch_replicator, requests_compressed, gzip]), + populate_db(Source, ?DOCS_COUNT), + replicate_with_options(Source, Target, [{<<"request_compression">>, <<"gzip">>}]), + compare_dbs(Source, Target), + After = couch_stats:sample([couch_replicator, requests_compressed, gzip]), + ?assert(After > Before) + after + config:delete("replicator", "compress_min_size", false) + end. + +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), + try + Before = couch_stats:sample([couch_replicator, requests_compressed, gzip]), + populate_db(Source, ?DOCS_COUNT), + replicate_with_options(Source, Target, [{<<"request_compression">>, <<"none">>}]), + compare_dbs(Source, Target), + After = couch_stats:sample([couch_replicator, requests_compressed, gzip]), + ?assertEqual(Before, After) + after + config:delete("replicator", "request_compression", false), + config:delete("replicator", "compress_min_size", false) + end. + +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). + +compare_dbs(Source, Target) -> + {ok, SourceInfo} = fabric:get_db_info(Source), + {ok, TargetInfo} = fabric:get_db_info(Target), + SourceDocCount = couch_util:get_value(doc_count, SourceInfo), + TargetDocCount = couch_util:get_value(doc_count, TargetInfo), + ?assertEqual(SourceDocCount, TargetDocCount), + ?assertEqual(?DOCS_COUNT, TargetDocCount).