diff --git a/Cargo.lock b/Cargo.lock index 3cae880092d6e..16b840b73b174 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3315,8 +3315,21 @@ checksum = "56254986775e3233ffa9c4d7d3faaf6d36a2c09d30b20687e9f88bc8bafc16c8" [[package]] name = "differential-dataflow" version = "0.25.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2677690d1dfbf4d8427eef2fdcf350ab0b258d3be20d01dd37dca07ca6f50bcc" +source = "git+https://github.com/DAlperin/differential-dataflow?rev=3659671246e0b5635a8128cc27d65f12ecf7dda5#3659671246e0b5635a8128cc27d65f12ecf7dda5" +dependencies = [ + "columnar", + "columnation", + "fnv", + "paste", + "serde", + "smallvec", + "timely 0.31.0 (git+https://github.com/TimelyDataflow/timely-dataflow)", +] + +[[package]] +name = "differential-dataflow" +version = "0.25.1" +source = "git+https://github.com/DAlperin/differential-dataflow?rev=70406df37670b1141904ce4395be0923c3e43125#70406df37670b1141904ce4395be0923c3e43125" dependencies = [ "columnar", "columnation", @@ -3324,7 +3337,7 @@ dependencies = [ "paste", "serde", "smallvec", - "timely", + "timely 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", ] [[package]] @@ -3333,9 +3346,9 @@ version = "0.25.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "52c9eee4acc936da45f096424617c91df981c619a92f3f8b4b2a324db7969b34" dependencies = [ - "differential-dataflow", + "differential-dataflow 0.25.1 (git+https://github.com/DAlperin/differential-dataflow?rev=70406df37670b1141904ce4395be0923c3e43125)", "serde", - "timely", + "timely 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", ] [[package]] @@ -6326,7 +6339,7 @@ dependencies = [ "datadriven", "dec", "derivative", - "differential-dataflow", + "differential-dataflow 0.25.1 (git+https://github.com/DAlperin/differential-dataflow?rev=70406df37670b1141904ce4395be0923c3e43125)", "enum-kinds", "fail", "futures", @@ -6398,7 +6411,7 @@ dependencies = [ "smallvec", "static_assertions", "thiserror 2.0.18", - "timely", + "timely 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", "tokio", "tokio-postgres", "tokio-stream", @@ -6422,7 +6435,7 @@ dependencies = [ "mz-storage-types", "serde", "serde_json", - "timely", + "timely 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", ] [[package]] @@ -6668,7 +6681,7 @@ dependencies = [ "bytesize", "chrono", "derivative", - "differential-dataflow", + "differential-dataflow 0.25.1 (git+https://github.com/DAlperin/differential-dataflow?rev=70406df37670b1141904ce4395be0923c3e43125)", "futures", "insta", "ipnet", @@ -6713,7 +6726,7 @@ dependencies = [ "sha2 0.10.9", "static_assertions", "thiserror 2.0.18", - "timely", + "timely 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", "tokio", "tracing", "uuid", @@ -6844,7 +6857,7 @@ version = "0.0.0" dependencies = [ "anyhow", "async-trait", - "differential-dataflow", + "differential-dataflow 0.25.1 (git+https://github.com/DAlperin/differential-dataflow?rev=70406df37670b1141904ce4395be0923c3e43125)", "futures", "lgalloc", "mz-cluster-client", @@ -6852,7 +6865,7 @@ dependencies = [ "mz-service", "rand 0.9.4", "regex", - "timely", + "timely 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", "tokio", "tracing", "tracing-subscriber", @@ -6884,7 +6897,7 @@ dependencies = [ "mz-dyncfg", "mz-ore", "mz-repr", - "timely", + "timely 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", "tokio", "tracing", ] @@ -6948,7 +6961,7 @@ dependencies = [ "mz-transform", "serde", "tempfile", - "timely", + "timely 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", "tokio", "tokio-stream", "tracing", @@ -6967,7 +6980,7 @@ dependencies = [ "core_affinity", "criterion", "dec", - "differential-dataflow", + "differential-dataflow 0.25.1 (git+https://github.com/DAlperin/differential-dataflow?rev=70406df37670b1141904ce4395be0923c3e43125)", "differential-dogs3", "fail", "futures", @@ -7002,7 +7015,7 @@ dependencies = [ "serde", "smallvec", "thiserror 2.0.18", - "timely", + "timely 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", "tokio", "tracing", "uuid", @@ -7018,7 +7031,7 @@ dependencies = [ "bytesize", "chrono", "derivative", - "differential-dataflow", + "differential-dataflow 0.25.1 (git+https://github.com/DAlperin/differential-dataflow?rev=70406df37670b1141904ce4395be0923c3e43125)", "futures", "mz-build-info", "mz-cluster-client", @@ -7039,7 +7052,7 @@ dependencies = [ "serde", "serde_json", "thiserror 2.0.18", - "timely", + "timely 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", "tokio", "tracing", "uuid", @@ -7052,7 +7065,7 @@ dependencies = [ "chrono", "chrono-tz", "columnar", - "differential-dataflow", + "differential-dataflow 0.25.1 (git+https://github.com/DAlperin/differential-dataflow?rev=70406df37670b1141904ce4395be0923c3e43125)", "fail", "itertools 0.14.0", "mz-dyncfg", @@ -7064,7 +7077,7 @@ dependencies = [ "serde", "serde-reflection", "serde_json", - "timely", + "timely 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", "tracing", ] @@ -7095,7 +7108,7 @@ dependencies = [ "regex", "serde", "serde_json", - "timely", + "timely 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", "tokio", "tracing", "uuid", @@ -7200,11 +7213,11 @@ dependencies = [ name = "mz-durable-cache" version = "0.0.0" dependencies = [ - "differential-dataflow", + "differential-dataflow 0.25.1 (git+https://github.com/DAlperin/differential-dataflow?rev=70406df37670b1141904ce4395be0923c3e43125)", "mz-ore", "mz-persist-client", "mz-persist-types", - "timely", + "timely 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", "tokio", "tracing", ] @@ -7377,7 +7390,7 @@ dependencies = [ "sysctl", "tempfile", "thiserror 2.0.18", - "timely", + "timely 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", "tokio", "tokio-metrics", "tokio-postgres", @@ -7588,7 +7601,7 @@ dependencies = [ "chrono", "clap", "criterion", - "differential-dataflow", + "differential-dataflow 0.25.1 (git+https://github.com/DAlperin/differential-dataflow?rev=70406df37670b1141904ce4395be0923c3e43125)", "itertools 0.14.0", "maplit", "mz-avro", @@ -7605,7 +7618,7 @@ dependencies = [ "seahash", "serde", "serde_json", - "timely", + "timely 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", "tokio", "tracing", "uuid", @@ -7936,7 +7949,7 @@ dependencies = [ "criterion", "ctor", "derivative", - "differential-dataflow", + "differential-dataflow 0.25.1 (git+https://github.com/DAlperin/differential-dataflow?rev=70406df37670b1141904ce4395be0923c3e43125)", "either", "futures", "hibitset", @@ -8023,7 +8036,7 @@ dependencies = [ "base64 0.22.1", "bytes", "deadpool-postgres", - "differential-dataflow", + "differential-dataflow 0.25.1 (git+https://github.com/DAlperin/differential-dataflow?rev=70406df37670b1141904ce4395be0923c3e43125)", "fail", "futures-util", "itertools 0.14.0", @@ -8049,7 +8062,7 @@ dependencies = [ "sha2 0.10.9", "tempfile", "time", - "timely", + "timely 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", "tokio", "tokio-postgres", "tracing", @@ -8072,7 +8085,7 @@ dependencies = [ "clap", "criterion", "datadriven", - "differential-dataflow", + "differential-dataflow 0.25.1 (git+https://github.com/DAlperin/differential-dataflow?rev=70406df37670b1141904ce4395be0923c3e43125)", "futures", "futures-task", "futures-util", @@ -8099,7 +8112,7 @@ dependencies = [ "serde", "serde_json", "tempfile", - "timely", + "timely 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", "tokio", "tokio-metrics", "tokio-stream", @@ -8409,7 +8422,7 @@ dependencies = [ "compact_bytes", "criterion", "dec", - "differential-dataflow", + "differential-dataflow 0.25.1 (git+https://github.com/DAlperin/differential-dataflow?rev=70406df37670b1141904ce4395be0923c3e43125)", "enum-kinds", "hex", "insta", @@ -8442,7 +8455,7 @@ dependencies = [ "static_assertions", "strsim", "thiserror 2.0.18", - "timely", + "timely 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", "tokio-postgres", "tracing", "tracing-core", @@ -8487,11 +8500,11 @@ dependencies = [ "ahash 0.8.12", "columnar", "columnation", - "differential-dataflow", + "differential-dataflow 0.25.1 (git+https://github.com/DAlperin/differential-dataflow?rev=70406df37670b1141904ce4395be0923c3e43125)", "mz-ore", "mz-repr", "mz-timely-util", - "timely", + "timely 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", ] [[package]] @@ -8754,7 +8767,7 @@ dependencies = [ "static_assertions", "thiserror 2.0.18", "tiberius", - "timely", + "timely 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", "tokio", "tokio-stream", "tokio-util", @@ -8852,7 +8865,7 @@ dependencies = [ "columnation", "crossbeam-channel", "dec", - "differential-dataflow", + "differential-dataflow 0.25.1 (git+https://github.com/DAlperin/differential-dataflow?rev=70406df37670b1141904ce4395be0923c3e43125)", "fail", "futures", "humantime", @@ -8910,7 +8923,7 @@ dependencies = [ "tempfile", "thiserror 2.0.18", "tiberius", - "timely", + "timely 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", "tokio", "tokio-postgres", "tokio-stream", @@ -8928,7 +8941,7 @@ dependencies = [ "aws-sdk-glue", "aws-smithy-mocks", "chrono", - "differential-dataflow", + "differential-dataflow 0.25.1 (git+https://github.com/DAlperin/differential-dataflow?rev=70406df37670b1141904ce4395be0923c3e43125)", "futures", "itertools 0.14.0", "mz-aws-glue-schema-registry", @@ -8952,7 +8965,7 @@ dependencies = [ "serde", "serde_json", "smallvec", - "timely", + "timely 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", "tokio", "tracing", "uuid", @@ -8966,7 +8979,7 @@ dependencies = [ "async-trait", "chrono", "derivative", - "differential-dataflow", + "differential-dataflow 0.25.1 (git+https://github.com/DAlperin/differential-dataflow?rev=70406df37670b1141904ce4395be0923c3e43125)", "futures", "itertools 0.14.0", "mz-build-info", @@ -8984,7 +8997,7 @@ dependencies = [ "mz-timely-util", "mz-txn-wal", "serde_json", - "timely", + "timely 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", "tokio", "tokio-stream", "tracing", @@ -9009,7 +9022,7 @@ dependencies = [ "columnation", "csv-async", "derivative", - "differential-dataflow", + "differential-dataflow 0.25.1 (git+https://github.com/DAlperin/differential-dataflow?rev=70406df37670b1141904ce4395be0923c3e43125)", "futures", "glob", "http 1.4.2", @@ -9034,7 +9047,7 @@ dependencies = [ "serde", "smallvec", "thiserror 2.0.18", - "timely", + "timely 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", "tokio", "tokio-util", "tracing", @@ -9062,7 +9075,7 @@ dependencies = [ "csv-core", "dec", "derivative", - "differential-dataflow", + "differential-dataflow 0.25.1 (git+https://github.com/DAlperin/differential-dataflow?rev=70406df37670b1141904ce4395be0923c3e43125)", "gcp_auth", "hex", "http 1.4.2", @@ -9114,7 +9127,7 @@ dependencies = [ "serde_json", "thiserror 2.0.18", "tiberius", - "timely", + "timely 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", "tokio", "tokio-postgres", "tracing", @@ -9224,7 +9237,8 @@ dependencies = [ "columnar", "columnation", "criterion", - "differential-dataflow", + "differential-dataflow 0.25.1 (git+https://github.com/DAlperin/differential-dataflow?rev=3659671246e0b5635a8128cc27d65f12ecf7dda5)", + "differential-dataflow 0.25.1 (git+https://github.com/DAlperin/differential-dataflow?rev=70406df37670b1141904ce4395be0923c3e43125)", "futures-util", "itertools 0.14.0", "lz4_flex 0.12.1", @@ -9235,7 +9249,8 @@ dependencies = [ "serde", "smallvec", "tempfile", - "timely", + "timely 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", + "timely 0.31.0 (git+https://github.com/TimelyDataflow/timely-dataflow)", "tokio", "tracing", "uuid", @@ -9290,7 +9305,7 @@ name = "mz-transform" version = "0.0.0" dependencies = [ "datadriven", - "differential-dataflow", + "differential-dataflow 0.25.1 (git+https://github.com/DAlperin/differential-dataflow?rev=70406df37670b1141904ce4395be0923c3e43125)", "enum-kinds", "fail", "itertools 0.14.0", @@ -9315,7 +9330,7 @@ version = "0.0.0" dependencies = [ "bytes", "crossbeam-channel", - "differential-dataflow", + "differential-dataflow 0.25.1 (git+https://github.com/DAlperin/differential-dataflow?rev=70406df37670b1141904ce4395be0923c3e43125)", "futures", "itertools 0.14.0", "mz-build-tools", @@ -9329,7 +9344,7 @@ dependencies = [ "prost-build", "rand 0.9.4", "serde", - "timely", + "timely 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", "tokio", "tracing", "uuid", @@ -10382,7 +10397,7 @@ dependencies = [ "axum", "bytes", "clap", - "differential-dataflow", + "differential-dataflow 0.25.1 (git+https://github.com/DAlperin/differential-dataflow?rev=70406df37670b1141904ce4395be0923c3e43125)", "futures", "humantime", "mz-dyncfg", @@ -10400,7 +10415,7 @@ dependencies = [ "prometheus", "serde", "serde_json", - "timely", + "timely 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", "tokio", "tracing", "uuid", @@ -11061,7 +11076,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "03da047801ff44bb6a4d407d4860c05fd70bb81714e6b2f3812603d5b145b042" dependencies = [ "heck", - "itertools 0.14.0", + "itertools 0.10.5", "log", "multimap", "petgraph", @@ -11082,7 +11097,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b570b25f7617e43d59005d0990ccb79e950a423952cea19671b7a876da390adf" dependencies = [ "anyhow", - "itertools 0.14.0", + "itertools 0.10.5", "proc-macro2", "quote", "syn 2.0.119", @@ -13483,10 +13498,28 @@ dependencies = [ "itertools 0.14.0", "serde", "smallvec", - "timely_bytes", - "timely_communication", - "timely_container", - "timely_logging", + "timely_bytes 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", + "timely_communication 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", + "timely_container 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", + "timely_logging 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", +] + +[[package]] +name = "timely" +version = "0.31.0" +source = "git+https://github.com/TimelyDataflow/timely-dataflow#c2c336c5a473f855b0fec02cd8b0c95d0abe0aa4" +dependencies = [ + "bincode", + "byteorder", + "columnar", + "columnation", + "itertools 0.14.0", + "serde", + "smallvec", + "timely_bytes 0.31.0 (git+https://github.com/TimelyDataflow/timely-dataflow)", + "timely_communication 0.31.0 (git+https://github.com/TimelyDataflow/timely-dataflow)", + "timely_container 0.31.0 (git+https://github.com/TimelyDataflow/timely-dataflow)", + "timely_logging 0.31.0 (git+https://github.com/TimelyDataflow/timely-dataflow)", ] [[package]] @@ -13495,6 +13528,11 @@ version = "0.31.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ad328384fad249cbe20f742d842ee9e3d2c6086bcde6277fb11bb154f21fb92e" +[[package]] +name = "timely_bytes" +version = "0.31.0" +source = "git+https://github.com/TimelyDataflow/timely-dataflow#c2c336c5a473f855b0fec02cd8b0c95d0abe0aa4" + [[package]] name = "timely_communication" version = "0.31.0" @@ -13505,9 +13543,22 @@ dependencies = [ "columnar", "getopts", "serde", - "timely_bytes", - "timely_container", - "timely_logging", + "timely_bytes 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", + "timely_container 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", + "timely_logging 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", +] + +[[package]] +name = "timely_communication" +version = "0.31.0" +source = "git+https://github.com/TimelyDataflow/timely-dataflow#c2c336c5a473f855b0fec02cd8b0c95d0abe0aa4" +dependencies = [ + "byteorder", + "columnar", + "serde", + "timely_bytes 0.31.0 (git+https://github.com/TimelyDataflow/timely-dataflow)", + "timely_container 0.31.0 (git+https://github.com/TimelyDataflow/timely-dataflow)", + "timely_logging 0.31.0 (git+https://github.com/TimelyDataflow/timely-dataflow)", ] [[package]] @@ -13516,13 +13567,26 @@ version = "0.31.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "de8c518a66e00ef8af9e9fae7e01b7d2628b970948f32855dfc89e503403ff5a" +[[package]] +name = "timely_container" +version = "0.31.0" +source = "git+https://github.com/TimelyDataflow/timely-dataflow#c2c336c5a473f855b0fec02cd8b0c95d0abe0aa4" + [[package]] name = "timely_logging" version = "0.31.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6ae9187f19aec95c932196b95c1981708090d9c35213ae266d7bac0f32432886" dependencies = [ - "timely_container", + "timely_container 0.31.0 (registry+https://github.com/rust-lang/crates.io-index)", +] + +[[package]] +name = "timely_logging" +version = "0.31.0" +source = "git+https://github.com/TimelyDataflow/timely-dataflow#c2c336c5a473f855b0fec02cd8b0c95d0abe0aa4" +dependencies = [ + "timely_container 0.31.0 (git+https://github.com/TimelyDataflow/timely-dataflow)", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index d6de0f6ce2336..af87fd24deb12 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -342,6 +342,7 @@ deadpool-postgres = "0.10.3" dec = "0.4.9" derivative = "2.2.0" differential-dataflow = "0.25.0" +differential-dataflow-next = { package = "differential-dataflow", git = "https://github.com/DAlperin/differential-dataflow", rev = "3659671246e0b5635a8128cc27d65f12ecf7dda5", default-features = false } differential-dogs3 = "0.25.0" digest = "0.10.7" dirs = "6.0.0" @@ -526,6 +527,7 @@ tikv-jemalloc-ctl = { version = "0.6", features = ["stats", "use_std"] } tikv-jemallocator = { version = "0.6", features = ["background_threads", "profiling", "stats", "unprefixed_malloc_on_supported_platforms"] } time = "0.3.17" timely = "0.31.0" +timely-next = { package = "timely", git = "https://github.com/TimelyDataflow/timely-dataflow", default-features = false } tokio = { version = "1.52.3", features = ["full", "test-util"] } tokio-metrics = "0.5.0" tokio-native-tls = "0.3.1" @@ -639,6 +641,9 @@ debug-assertions = true # merged), after which point it becomes impossible to build that historical # version of Materialize. [patch.crates-io] +# Synchronous spine exertion bounded by inserted updates and frontier advances. +# Branch upsert-funded-exertion-0.25.1 on top of the 0.25.1 release. +differential-dataflow = { git = "https://github.com/DAlperin/differential-dataflow", rev = "70406df37670b1141904ce4395be0923c3e43125" } # Waiting on https://github.com/sfackler/rust-postgres/pull/752. postgres = { git = "https://github.com/MaterializeInc/rust-postgres" } tokio-postgres = { git = "https://github.com/MaterializeInc/rust-postgres" } diff --git a/deny.toml b/deny.toml index ec3704b00ef8a..0c642a679d0b9 100644 --- a/deny.toml +++ b/deny.toml @@ -16,6 +16,14 @@ multiple-versions = "deny" # Materialize-maintained fork that avoids the duplicated transitive # dependencies. skip = [ + # The DNM async-merge prototype consumes two DD/Timely APIs. + { name = "differential-dataflow", version = "=0.25.1" }, + { name = "timely", version = "=0.31.0" }, + { name = "timely_bytes", version = "=0.31.0" }, + { name = "timely_communication", version = "=0.31.0" }, + { name = "timely_container", version = "=0.31.0" }, + { name = "timely_logging", version = "=0.31.0" }, + # arrayvec had a significant API change in 0.7 { name = "arrayvec", version = "0.5.2" }, # One-time exception for base64 due to its prevalence in the crate graph. @@ -367,6 +375,11 @@ license-files = [ [sources] unknown-git = "deny" unknown-registry = "deny" +# DNM-only exceptions for the async-merge prototype, pinned in Cargo.lock. +allow-git = [ + "https://github.com/DAlperin/differential-dataflow", + "https://github.com/TimelyDataflow/timely-dataflow", +] # Do not allow non-MaterializeInc Git repositories here! Git repositories must # be owned by the MaterializeInc organization so that maintainership is shared # amongst Materialize employees and so that historical versions of Materialize diff --git a/doc/developer/design/20260908_out_of_core_operators.md b/doc/developer/design/20260908_out_of_core_operators.md new file mode 100644 index 0000000000000..3c91c8c658852 --- /dev/null +++ b/doc/developer/design/20260908_out_of_core_operators.md @@ -0,0 +1,572 @@ +# A shared model for out-of-core operators + +- Status: discussion draft, September 8, 2026. Interfaces and rollout gates below are proposals. +- Related: [Buffer-managed dataflow state](20260610_buffer_managed_state.md). +- Related experiment: [upsert hydration optimizations, PR #38719](https://github.com/MaterializeInc/materialize/pull/38719). + +## Proposal in brief + +Separate operator state into immutable payload storage and a spillable index of +keys, handles, timestamps, and differences. Execute operators through resumable, +budgeted work units that request payloads only when their semantics require them. +Reuse Differential's time and difference machinery where possible, including the +boundary demonstrated by `int_proxy`. + +Upsert is the first consumer. Equijoin is the second design and implementation +check. Neither storage nor scheduling should know about source offsets, +latest-value selection, join predicates, or a particular reduction function. + +The central hypothesis is that moving compact metadata through repeated merges, +while keeping payload blocks independently owned, reduces byte amplification. +Whether this improves elapsed time depends on the additional lookup, equality, +fragmentation, and payload-read costs. We will measure those costs separately. + +## The problem + +Our chunked state paths can repeatedly serialize, copy, compress, and decompress +wide rows while sorting or merging updates. Much of that work is needed to move +state through the representation, even when the operation only needs a key, +timestamp, or source ordering field to make its decision. + +The existing pool controls chunk residency. It does not by itself separate payload +lifetime from merge lifetime. Smaller chunks, tighter admission, direct compressed +output, and offloaded reads improve the current representation. A shared operator +model should also let a merge keep a payload reference instead of rebuilding the +payload at every generation. + +This is broader than upsert hydration. Joins need to find matching groups and +materialize selected pairs. Reductions need to replay histories and reconcile +outputs. All need progress tracking, bounded staging, storage ownership, and an +execution path that can wait for unavailable data without blocking a worker. + +## Success criteria + +- Upsert and equijoin share payload storage, ownership, memory accounting, and + read scheduling. Their semantic rules remain in their operator implementations. +- A metadata merge can retain live payloads without decoding or rewriting them. + Repacking payload blocks is separately scheduled and measured. +- Payloads, indexes, handle translation, ownership metadata, staging, and queued + work have explicit memory charges. Increasing state beyond the budget does not + introduce an unbounded resident table of row handles or block descriptors. +- Supported operators make progress under a small budget, including with skew, + slow reads, output backpressure, cancellation, and concurrent consumers. +- Results and progress agree with existing operators under arbitrary valid + batching, retractions, timestamp order, compaction, and restarts. +- Deep-state workloads improve without an unacceptable resident-state penalty. + Proposed acceptance thresholds appear under measurement, rather than assuming + that fewer copied bytes guarantee faster queries. + +## Scope + +This proposal covers recreatable local operator state. Persist remains the durable +source of recovery. Local handles do not become part of persisted records or the +wire protocol. Crossing a process boundary materializes or explicitly transfers +owned data rather than sending a process-local handle. + +We will preserve existing in-memory implementations during evaluation. Converting +all operators, implementing a new durable store, and selecting a particular kernel +async I/O API are outside the first delivery. General reduction remains a design +requirement, but arbitrary reduction callbacks do not automatically acquire a +bounded-memory implementation through this interface. + +## What int-proxy contributes + +Differential's `int_proxy` tactics exchange consolidated presentations of: + +```text +((key_hash: u64, value_id: u64), time, diff) +``` + +The tactic performs the time and difference computation. The backend interprets +values and constructs outputs. A join returns matched IDs with joined times and +multiplied differences. A reduce asks the backend for corrections at the times +that require reconciliation. + +The two integers have different contracts. `key_hash` partitions independent +work, but collisions are allowed. `value_id` identifies data within the backend's +presentation. Distinct logical data in the same hash group must remain distinct, +and equal data must be presented consistently wherever cancellation is required. +Input and output presentations can have separate ID namespaces. + +The local reference reduce backend assigns IDs per window and resolves them +through row vectors. It still clones rows and emits ordinary row batches. The +proxy interface therefore demonstrates a separation of responsibilities, not a +ready-made payload store or an out-of-core guarantee. + +Two constraints matter for this design: + +1. The callbacks are synchronous. A cold read cannot simply become an `await` + inside the existing tactic protocol. +2. The reduce window contract requires a complete key-hash group in the window + that first reports it. A window target does not bound a single enormous group. + Its novel time support must also survive consolidation: even an update that + cancels against history can introduce a time requiring reconciliation. + +Reference code is in the sibling Differential repository under +`differential-dataflow/src/operators/int_proxy/{mod,reduce,join,vec_backend}.rs`. +This draft describes that local implementation. Materialize currently locks +Differential 0.25.1, whose proxy API must be reconciled with the local work before +integration. Names and signatures here are illustrative, not a compatibility claim. + +## Architecture and ownership boundaries + +```mermaid +flowchart TD + A[Operator semantics: upsert, join, reduce] --> B[Resumable execution and time/diff logic] + B --> C[Spillable indexes and immutable batches] + B --> D[Budgeted payload reads] + C --> E[Payload store and ownership manifests] + D --> E + E --> F[Pool and extent storage] + G[Shared memory and I/O admission] -.-> B + G -.-> C + G -.-> D + G -.-> E +``` + +### Payload storage + +Store immutable row blocks independently of index batches. A block includes its +record boundaries and encoding information. The store supports batched resolution +of row handles, allowing requests to be grouped by block and decoded together. +The first implementation should reuse the pool's codecs and extent machinery. + +Reads initially return owned, budgeted decoded buffers. References into a decoded +buffer are scoped to a read lease. This follows the pool's current copy-out model +and avoids requiring a redesign around borrowed pointers into evictable slots. +A lease retains both the bytes' memory charge and ownership of the source needed +by an unfinished read. Callers cannot retain a naked row reference after the +lease expires. + +Block size is a policy choice, separate from index chunk size. Smaller blocks +reduce sparse-read amplification but add metadata and can weaken compression. +Benchmark several sizes rather than baking the upsert workload's preferred size +into the interface. Oversized individual rows require a charged exceptional path +or streaming support, never an unaccounted allocation. + +### Identity: handles need not be integers + +Use distinct types for storage identity and presentation identity: + +| Identity | Meaning | Lifetime | +| --- | --- | --- | +| `RowHandle` | Locates a stored row within a store namespace | While owning batches or read operations retain it | +| `ProxyId` | Identifies semantic data in an operator presentation | Until that presentation's work and output translation finish | +| Group key | Defines independent semantic work | Defined by the operator and its index | + +A candidate `RowHandle` representation is `(block_id, row_slot, generation)`. +Moving a block between resident memory and an extent changes its location, not +its identity. Generation checks prevent stale handles from resolving to reused +storage. The store namespace is supplied by the owning batch or execution context. + +The current proxy bridge requires `u64`, but the design does not require every +identity to be a `u64`. A backend can translate storage handles to dense integer +IDs for a bounded window. We should keep that adapter before generalizing the +upstream tactic's types unless two consumers demonstrate a need for the latter. +A handle's numeric order does not imply row order. + +Physical identity is also not semantic equality. Two independently ingested equal +rows can have different handles. Exact consolidation needs an equality mechanism: +for example, a fingerprint to find candidates followed by byte comparison, with +consistent proxy IDs assigned to equal data during presentation. Fingerprints +alone are insufficient. Such lookup state must be budgeted and spillable too. + +Do not require global interning of all live rows. Window-local canonicalization +and merge-time equality resolution are the starting point. Operators can exploit +stronger local knowledge, such as an already-owned output being reused, without +making that knowledge a storage requirement. + +### Indexes and batches + +An index entry carries the fields required for navigation and operator decisions, +plus references to payloads. Field projection belongs to an index layout or +operator backend. The storage layer does not assume a fixed tuple containing +upsert-specific ordering information. + +Keys can themselves be wide. A hash or compact prefix may identify candidate +groups, but exact key comparison can require payload reads. The execution protocol +must permit these reads during navigation and consolidation, not just at final +output. Likewise, a reducer that requires value ordering cannot sort opaque +handles and assume the order matches values. + +Index chunks, fences, equality indexes, and manifests must support external +storage. Resident roots and caches have byte budgets. Existing merge machinery +can be reused where it operates on the chosen compact representation. Comparator +or equality paths requiring cold data need an explicit suspend/resume boundary. + +### Ownership and reclamation + +A `RowHandle` is a locator, not an owning reference. An immutable batch manifest +owns the payload blocks its entries reference. Active reads and unpublished output +builders acquire ownership too. Prefer ownership per block or segment over an +atomic reference count for every row copy. + +Publishing an output batch must establish its manifest's ownership before input +ownership is released. Cancellation discards unpublished output and releases its +ownership. Dropping a batch releases references incrementally, with pending +release work charged and able to yield. A block is reclaimable only when no +published batch, builder, or active read owns it. A timestamp frontier alone is +not a reclamation proof. + +Coarse ownership can retain mostly dead blocks. Repacking live rows is a separate, +budgeted maintenance operation that creates new blocks and rewrites affected index +references. Old batches and active readers continue to own old blocks until they +retire. Do not introduce an unbounded per-row forwarding map to hide relocation. +Measure the temporary double ownership and charge the rewrite's scratch space. + +The ownership manifest and block directory are substantial parts of the work. +Retaining one always-resident `Arc` per live block merely moves the +scaling limit. A production design needs paged directory/ownership metadata, +bounded resident roots, and a way to retire metadata without loading an entire +batch manifest. The current pool API may need extensions at this boundary. + +## Execution, progress, and budgets + +The execution unit is a continuation with explicit input ownership, progress +holds, and a bounded working set. Its conceptual protocol is: + +```text +advance(work_budget) + -> need_reads(read_set, continuation) + | produced(output_batch, continuation) + | yield(continuation) + | complete + | error +``` + +These are protocol outcomes, not proposed Rust signatures. Metadata and payload +reads use the same admission rules. A ready-data step is synchronous and bounded +by bytes or work, with a separate fairness limit so a long CPU-only step yields. +Output backpressure suspends execution without accumulating unlimited ready output. + +The driver reserves decode/output memory and I/O capacity before issuing a read. +Completed buffers remain charged until consumed, including while the worker is +busy elsewhere. A limit on running reads alone is insufficient. Several operators +share process-level admission, with per-consumer fairness and an allowance for +work that releases memory. An operator must not hold the whole budget while +waiting for an additional allocation needed to make progress. The prototype must +exercise this deadlock case and establish a reservation/release discipline. + +No pool locks, borrowed cursors into mutable state, or uncharged buffers may cross +a suspension. Cancellation retains a running job's lease and permit until that +job actually finishes, then releases them even if its consumer disappeared. +Storage errors invalidate the work and are surfaced through the operator's error +or restart path. Partially built output is not published as a completed batch. + +Suspension must preserve Differential's time semantics. Work retains capabilities +or equivalent holds for every timestamp at which it can still emit. A yield or +pending read must not advance the output frontier. Publication and continuation +updates must avoid duplicate output when work resumes. Partial-order timestamps, +compaction, and novel time support remain the responsibility of the time/diff +machinery and its adapter, rather than the payload store. + +Time/diff bookkeeping is part of the budget too. Replay histories, pending-time +schedules, and retained presentations must be spillable or have an explicit +admission bound. The tactic's existing resident pending-time map cannot be treated +as free metadata. Any bound must preserve required times and holds, using +backpressure or external state rather than dropping work. This is a dependency of +the resumable tactic design, not something the payload store can solve alone. + +There are two integration options: prepare a bounded, fully resident presentation +before invoking a synchronous tactic, or make tactic execution resumable. The +first is useful for a prototype, but cannot provide the complete contract for +large groups or payload-dependent comparisons. The target design requires a +resumable path. We should upstream the smallest suitable execution boundary +rather than duplicate Differential's time logic in each backend. + +The initial I/O implementation can offload swap faults and decoding to a bounded +blocking executor. It does not guarantee fault-free worker execution, and it is +not native asynchronous file I/O. The same request/continuation contract should +permit a later file-extent implementation without changing operator semantics. + +## Large groups and operator contracts + +A hot key can exceed any window budget. Different operators need different +strategies, and declaring an arbitrary split into independent windows is incorrect. + +| Consumer | Small representation | Payload access | Large-group strategy | +| --- | --- | --- | --- | +| Upsert | Key, time, source order, optional row handle | Equality and old/new rows needed for output | Stream latest-update selection while preserving the existing feedback and frontier rules | +| Equijoin | Key projection/hash, row handles, time, diff | Exact key checks, predicates, output construction | Block both sides, retain resumable history, and stream matches under output backpressure | +| Threshold or aggregatable reduce | Group identity, value identity or sufficient aggregate, time, diff | Depends on equality and aggregate semantics | Use valid partial state or spillable histories | +| General reduce | Group and value references with histories | Potentially all values in semantic order | Require a streaming callback or external intermediate state; an arbitrary slice callback has no bounded-memory guarantee | + +Upsert's source stash can select by source order without reading every row. +Its feedback state still requires exact row identity for retractions and +consolidation. Equality reconciliation after persist readback must not assume +that a deserialized row preserves its former local handle. + +Equijoin must enumerate all valid matches even when output greatly exceeds input. +That cost cannot be removed by proxies. Blocking may require rereading one side, +so we must measure actual read amplification and scheduler fairness. + +These two consumers should use the same lower layers. If implementing equijoin +requires a second payload store or ownership mechanism, the proposed abstraction +has not met its purpose. + +## Minimal viable prototype and delivery + +This document precedes a prototype so we can agree on the ownership and execution +boundaries before implementing them. No performance gains from this representation +have been measured yet. + +1. Build a small shared store/index experiment with a deliberately tiny budget, + duplicate logical rows at different handles, asynchronous reads, and cancellation. + Exercise both an upsert-like selector and an equijoin over it. Keep the interfaces + provisional until both work. A fully resident metadata directory can be used to + establish mechanics, but cannot satisfy the out-of-core acceptance gate. +2. Implement spillable metadata and manifests, exact equality reconciliation, + reclamation, large-group continuations, and publication/progress invariants. + Demonstrate budget and liveness properties independently of hydration speed. +3. Integrate upsert while preserving its existing time and feedback semantics. + Compare against v1 and the measured optimization stack. Keep persisted and wire + representations unchanged so a fresh replica can select either implementation. +4. Integrate equijoin using the same store, admission, and ownership APIs. Validate + skew and output backpressure. Use a reduction consumer to settle remaining + value-order and large-group interface questions before claiming broad support. +5. Roll out per consumer behind flags. Production defaults off, test/CI defaults + on for supported paths. Check small-state overhead and deep-state behavior before + increasing exposure. Rollback recreates local state from persist rather than + converting live handles between formats. + +## Correctness and measurement + +Correctness testing compares output updates and progress with existing operators, +including product timestamps with incomparable elements. Vary input batch order, +negative differences, frontier advancement, compaction, empty inputs, deletion, +and duplicate rows. Force hash collisions and separate physical copies of equal +data. Preserve the case where novel input cancels against history but still +introduces a time requiring reconciliation. + +Storage and scheduler tests cover eviction during reads, cancellation before and +after submission, delayed completion, stale handle generations, output publication +before input release, repacking while old batches remain live, scratch exhaustion, +and injected read failures. Run with budgets smaller than a group and with two +competing consumers. Verify both progress and eventual release of charges. +Restart tests reconstruct state from persist and verify that no local handles leak +into durable or exchanged data. + +Use the spec-sheet persisted-seed workflow for upsert hydration: load once, stop +the writer, and recreate a replica for each variant. Record the exact image, +configuration, seed shape, worker count, CPU, memory, and usable scratch/swap. +Compare the current optimized representation with the proposed one on the same +build where possible. Keep a v1 control, but do not make v1 parity the only success +criterion for a shared framework. + +Measure equijoin separately with fixed input/output cardinalities, selective and +dense matches, skew, and wide join keys. Add steady-state update/retraction runs +and concurrent consumers. The first upsert screening replica has one worker and +half a CPU; it cannot establish production concurrency or scheduling behavior. + +Sweep payload width, actual compressibility, state-to-budget ratio, update density, +key skew, worker/core ratio, and payload block size in stages rather than running +an undirected Cartesian product. Include resident-state controls and poorly +compressible data. Randomize or counterbalance variant order, retain failed runs, +and repeat enough trials to report variation rather than one timing. + +Collect elapsed time, CPU, worker scheduling delay, time waiting for reads, decode +time, peak RSS, swap, and all memory-ledger components. Count metadata merge bytes, +payload bytes copied/decoded/rewritten, equality lookups and verification reads, +block cache hits, dead-byte retention, reclamation work, and output bytes. Separate +extent creation from physical device traffic. Account for sampling gaps and missing +metrics explicitly. Treat restarts as failed trials, not successful slow runs. + +Proposed gates for discussion: + +- Correctness and bounded-memory/liveness checks pass for both initial consumers. +- Supported deep-state workloads finish under a fixed budget without growing an + unaccounted directory, pending-read queue, or output queue. +- On wide-row, deep-state cases with comparable output, reduce payload bytes + processed by state maintenance by at least 2x versus the optimized chunk path, + and demonstrate an elapsed-time improvement outside run-to-run variation. +- Target no more than 5% resident-state slowdown for each consumer. If overhead + exceeds that threshold, investigate layout/admission costs or keep an explicit + small-state path. These thresholds are proposed acceptance criteria, not forecasts. + +## Local mechanics prototype + +The prototype lives in [`mz_timely_util::out_of_core`](../../../src/timely-util/src/out_of_core.rs). +It provides a shared payload store, immutable block manifests, resident metadata +batches, and whole-step read admission. A read request owns its blocks before it +can suspend. A blocking read job owns its byte and job permits through +cancellation, and the returned lease retains the byte reservation until its +borrowed output is consumed. + +Two consumers share these interfaces: maximum-order selection and equijoin. +Selection rebuilds metadata while reusing existing payload blocks. The resumable +join cursor emits one pair per step and uses the same store to acquire both +payloads atomically. String and composite keys work without integer encoding. + +The standalone proxy adapters and live Differential harness are retained in the +local prototype workspace. The benchmark image includes the shared storage and +native columnar upsert integration below, using Materialize's existing Timely +and Differential dependencies. It has no sibling-checkout dependency. + +## Live dataflow setup + +`out_of_core::live` in the standalone prototype workspace installs +operators into a running Timely dataflow using the local Differential checkout: + +1. `store_payloads` exchanges serialized input rows by key, then packs their + bytes into the worker's payload store. It emits `Stored` values whose + ordering and equality use a caller-defined logical identity. Each value owns + its payload block. Handles never cross workers. +2. `arrange_payloads` builds a Differential chunk spine over key hashes, exact + keys, and stored values. Trace consolidation and merging clone metadata and + ownership references, without copying payload bytes. +3. `latest` runs `ProxyReduceTactic` through `reduce_with_tactic`. `join` runs + `ProxyJoinTactic` through `join_with_tactic`. These drivers maintain input and + output history, compaction frontiers, and capabilities across successive + input batches. The backends resolve exact identities and hash collisions. +4. `fetch` requests the complete read set for an output record, polls its future + on the Timely worker, and uses a synchronous activator to resume when I/O + completes. Blocking pool reads run through the Tokio executor. Each future + holds an output capability and its stored values until publication. Incoming + antichain stamps are retained as capability sets before deriving a capability + at each record's time. + +The integration harness runs selection and join together on two workers and +compares both outputs against standard Differential operators. Inputs advance +through multiple epochs and are probed before sending subsequent epochs. It also +covers product timestamps, two-element message stamps, error propagation, a +blocked read that leaves another dataflow schedulable, and dropping a dataflow +while its blocking read still owns admission. + +```sh +MZ_DEV_BUILD_SHA=f0e632d1 cargo +1.97.1 nextest run \ + -p mz-timely-util --features out-of-core-prototype --test out_of_core_live + +MZ_OOC_EPOCHS=1000 MZ_DEV_BUILD_SHA=f0e632d1 cargo +1.97.1 nextest run \ + -p mz-timely-util --features out-of-core-prototype --test out_of_core_live \ + -E 'test(live_proxy_pipeline_matches_reference_on_two_workers)' --no-capture +``` + +A local 1,000-epoch run matched both reference outputs. With a 1,024-byte decoded +budget per worker, observed peaks were 1,024 and 992 bytes. Live payload blocks +sampled after each completed epoch peaked at 10 and 5, and both workers released +all blocks on closure. This is a functional small-state run over compressed +extents, not a hydration or disk-throughput benchmark. + +The generic live harness remains separate from Materialize's source and compute +renderers. The native upsert integration below uses the same storage and ownership +model with Materialize's current Differential dependency and existing feedback +protocol. + +The remaining shared runtime work includes spillable metadata and ownership +indexes, exact row equality where the caller has no logical identity, resumable +window presentation and history loading, decoded-block reuse, and admission covering +queued metadata and output. Production wiring also needs one compatible Timely +and Differential dependency set, row/error codecs, operator shutdown integration, +metrics, configuration defaults, and restart/rehydration validation. + +### Resumable join matching + +The local proxy join retains direct-cross positions or bilinear replay histories +across calls. Each prepared work unit clones its configured backend, giving it a +private ID interpretation table. A window can produce multiple `cross` calls, +so that table remains valid until the next window or work-unit drop. + +`JoinWork::step` returns output, yield, or done. The live driver retains its +capabilities on yield and reactivates the operator, including when exact-key +filtering discards every candidate. Matching defaults to 4,096 candidate matches +and 4,096 replay transitions per quantum. The compatibility iterator consumes +yields internally and provides no scheduling guarantee. + +Window presentation, history loading, and consolidation remain synchronous and +can exceed a quantum. The limits bound staged matches and matching transitions, +not the whole activation, input memory, or arbitrary backend output expansion. +A live test checks a 10,000-pair hot key through pool-backed payload fetches. + + +### Native columnar upsert integration + +`enable_upsert_payload_stash` selects a third upsert-v2 state representation. It +is off in production and on in the mzcompose test parameter defaults. It takes +precedence over `enable_upsert_chunked_stash` when upsert-v2 is enabled. + +`columnar::payload::PayloadChunk` wraps a normal columnar metadata chunk and a +manifest. The existing Differential chunk batcher and spine handle merging, +advancement, sealing, and compaction. Metadata can spill through the existing +columnar path. Each rewrite retains exactly the referenced payload blocks, +without decoding payloads. The generic bulk-probe interface returns metadata +with the ownership needed to resolve its locators. + +The source stash stores keys, times, source offsets, tombstones, and row handles. +Offset selection is the existing `UpsertDiff` semigroup. Payloads are published +after the initial chunker selects its winners, in metadata order. Ineligible updates retain +their handles when re-stashed. The feedback arrangement stores keys and exact +payload identities with ordinary additive diffs. An operator-local weak equality +index uses fingerprints to find candidates, then checks their bytes before +reusing a live handle. It owns no blocks and sweeps dead entries incrementally. +This identity policy is specific to this consumer, rather than required by the +chunk abstraction. Its resident index cost needs measurement. + +Source and feedback share one payload store: 2 MiB blocks, a 16 MiB decoded-read +budget, and two blocking read jobs. The drain retains at most two decoded blocks +for reuse across windows. Each 1,024-record window groups emitted updates by +payload locator before decoding, avoiding old/new block alternation. Feedback +canonicalization retains one decoded block. Payloads larger +than a block remain inline, preserving supported row sizes at the cost of the +original row-copying behavior for those rows. Encoded error rows follow the same +path as successful values. The persist format is unchanged. + +The upsert driver still owns resume filtering, persist eligibility, frontiers, +metrics, and shutdown. This does not replace that protocol with a generic latest +reduction, and does not install the local proxy join driver into compute. It +requires no Timely or Differential dependency migration. + +Limits to measure include the resident equality index and manifests, metadata +reads while pruning ownership, synchronous metadata merges, batching scratch, +and payload block retention when only a few rows remain live. Decoded admission +is shared by the source and feedback operators of one upsert dataflow, not yet +across every operator in the process. The process pool still governs residency. + +The hydration comparison adds `v2_payload` alongside `v1` and `v2_all`, using +identical persisted state and fresh replicas. The candidate image must contain +the new flag. The existing overnight image predates this integration. + + +## Alternatives + +**Continue improving combined row chunks.** This has the smallest implementation +cost and preserves write elision for short-lived chunks. It remains the baseline. +It is sufficient if measurements show that byte amplification is no longer the +limiting cost, and avoids the proposed equality and ownership overhead. + +**Adopt the existing int-proxy reference backend directly.** This validates tactic +integration, but does not eliminate row copying or provide bounded cold reads, +metadata, or large-group execution. It is a useful reference, not the target store. + +**Globally intern every row to an integer.** Equality becomes cheap after lookup, +but the interning index and ownership can dominate memory and I/O. Requiring it +also couples every consumer to a global identity policy. Prefer bounded exact +canonicalization, with stronger interning optional when measurements justify it. + +**Use one key/value store for all operator state.** This can supply storage and +lookup, but still needs integration with immutable batch sharing, time histories, +progress, and output backpressure. It may be a backend option rather than the +operator interface itself. + +**Make native async I/O the first project.** It can improve read scheduling, but +preserves unnecessary bytes read, written, and decoded. Separate the operator's +resumable contract from the extent implementation so both improvements can be +measured independently. + +## Questions for discussion + +1. **Execution boundary:** should Differential expose resumable tactics, or a + smaller replay primitive that a Materialize async driver composes? Preparing + resident windows is a useful first experiment, not a complete large-group answer. +2. **Ownership granularity:** can block/segment manifests and paged ownership + metadata provide acceptable dead-byte retention, or do we need finer liveness + information? What is the minimum pool API extension required? +3. **Equality:** which consumers can canonicalize within a merge or work window, + and which need a longer-lived spillable equality index? How much cold comparison + traffic does each approach introduce? +4. **Large groups:** which reduction contracts are worth supporting initially, + and how should operators declare their streaming or external-state requirements? +5. **Delivery gate:** do we agree that equijoin must exercise the shared runtime + before these APIs are considered stable, even though upsert ships first? +6. **Performance gate:** are the proposed 2x byte-reduction and 5% resident-overhead + targets appropriate, and which workloads should decide whether the additional + storage machinery is justified? diff --git a/doc/user/data/metrics.yml b/doc/user/data/metrics.yml index d42be02e4edd0..c105e251d2e38 100644 --- a/doc/user/data/metrics.yml +++ b/doc/user/data/metrics.yml @@ -556,6 +556,22 @@ metrics: help: Evicted chunks re-admitted by an admitting read stealing the slot of a clean backed victim of the same size class. source: src/timely-util/src/pool_config/metrics.rs visibility: internal +- name: mz_column_pool_async_reads_in_flight + help: Submitted pool reads that have not released their concurrency permit. + source: src/timely-util/src/pool_config/metrics.rs + visibility: internal +- name: mz_column_pool_async_reads_total + help: Pool reads submitted to the blocking executor. + source: src/timely-util/src/pool_config/metrics.rs + visibility: internal +- name: mz_column_pool_cold_inserts_total + help: Chunks inserted directly into an extent by caller request. + source: src/timely-util/src/pool_config/metrics.rs + visibility: internal +- name: mz_column_pool_direct_extent_inserts_total + help: Inserts written directly to an extent because resident admission was full. + source: src/timely-util/src/pool_config/metrics.rs + visibility: internal - name: mz_column_pool_eager_backs_total help: Chunks eagerly compressed to compressed-but-resident by idle spill threads; their later eviction is a pure page release. source: src/timely-util/src/pool_config/metrics.rs diff --git a/misc/python/materialize/mzcompose/__init__.py b/misc/python/materialize/mzcompose/__init__.py index a2e156a9e627b..303f4dc6d359c 100644 --- a/misc/python/materialize/mzcompose/__init__.py +++ b/misc/python/materialize/mzcompose/__init__.py @@ -364,6 +364,26 @@ def get_variable_system_parameters( "true", ["true", "false"], ), + VariableSystemParameter( + "enable_column_chunk_direct_compressed_output", + "true", + ["true", "false"], + ), + VariableSystemParameter( + "enable_upsert_payload_stash", + "true", + ["true", "false"], + ), + VariableSystemParameter( + "enable_upsert_async_merges", + "true", + ["true", "false"], + ), + VariableSystemParameter( + "enable_upsert_async_reads", + "true", + ["true", "false"], + ), VariableSystemParameter( "enable_upsert_v2", "false", diff --git a/misc/python/materialize/parallel_workload/action.py b/misc/python/materialize/parallel_workload/action.py index 6a008fc235834..00b06ea163b1a 100644 --- a/misc/python/materialize/parallel_workload/action.py +++ b/misc/python/materialize/parallel_workload/action.py @@ -3179,6 +3179,11 @@ def __init__( ] self.flags_with_values["enable_upsert_paged_spill"] = BOOLEAN_FLAG_VALUES self.flags_with_values["enable_upsert_chunked_stash"] = BOOLEAN_FLAG_VALUES + self.flags_with_values["enable_upsert_payload_stash"] = BOOLEAN_FLAG_VALUES + self.flags_with_values["enable_upsert_async_reads"] = BOOLEAN_FLAG_VALUES + self.flags_with_values["enable_column_chunk_direct_compressed_output"] = ( + BOOLEAN_FLAG_VALUES + ) self.flags_with_values["column_chunk_compress_min_depth"] = [ "0", # compress every spilled body "1", # the default: fresh chunks store uncompressed diff --git a/src/cluster/src/client.rs b/src/cluster/src/client.rs index 0008e72b0d7fe..b335f62c15229 100644 --- a/src/cluster/src/client.rs +++ b/src/cluster/src/client.rs @@ -10,6 +10,7 @@ //! An interactive cluster server. use std::fmt; +use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::{Arc, Mutex}; use std::thread::Thread; @@ -18,6 +19,8 @@ use async_trait::async_trait; use differential_dataflow::trace::ExertionLogic; use futures::future; use mz_cluster_client::client::{TimelyConfig, TryIntoProtocolNonce}; +use mz_ore::metric; +use mz_ore::metrics::{ComputedUIntGauge, MetricsRegistry}; use mz_service::client::{GenericClient, Partitionable, Partitioned}; use mz_service::local::LocalClient; use timely::WorkerConfig; @@ -225,6 +228,7 @@ pub trait ClusterSpec: Clone + Send + Sync + 'static { // largest to the smallest layer. let arc: ExertionLogic = Arc::new(move |layers| { + EXERT_POLICY_CALLS.fetch_add(1, Ordering::Relaxed); let mut prop = config.arrangement_exert_proportionality; // Layers are ordered from largest to smallest. @@ -238,12 +242,14 @@ pub trait ClusterSpec: Clone + Send + Sync + 'static { for (_idx, count, len) in layers { if count > 1 { // Found an in-progress merge that we should continue. + EXERT_POLICY_MERGE_GRANTS.fetch_add(1, Ordering::Relaxed); return merge_effort; } if !first && prop > 0 && len > 0 { // Found a non-empty batch within `arrangement_exert_proportionality` of // the largest one. + EXERT_POLICY_CONSOLIDATION_GRANTS.fetch_add(1, Ordering::Relaxed); return merge_effort; } @@ -284,6 +290,29 @@ pub trait ClusterSpec: Clone + Send + Sync + 'static { } } +// The exertion policy runs on every arrangement's scheduling turn, so these +// process-wide counts describe how often traces of every implementation were +// offered optional effort, and for which reason. +static EXERT_POLICY_CALLS: AtomicU64 = AtomicU64::new(0); +static EXERT_POLICY_MERGE_GRANTS: AtomicU64 = AtomicU64::new(0); +static EXERT_POLICY_CONSOLIDATION_GRANTS: AtomicU64 = AtomicU64::new(0); + +/// Register gauges for the arrangement exertion policy's decisions. +pub fn register_exert_policy_metrics(registry: &MetricsRegistry) { + let _: ComputedUIntGauge = registry.register_computed_gauge( + metric!(name: "mz_arrangement_exert_policy_calls_total", help: "Arrangement maintenance turns that consulted the exertion policy."), + || EXERT_POLICY_CALLS.load(Ordering::Relaxed), + ); + let _: ComputedUIntGauge = registry.register_computed_gauge( + metric!(name: "mz_arrangement_exert_policy_grants_total", help: "Exertion policy grants, by reason.", const_labels: {"reason" => "active_merge"}), + || EXERT_POLICY_MERGE_GRANTS.load(Ordering::Relaxed), + ); + let _: ComputedUIntGauge = registry.register_computed_gauge( + metric!(name: "mz_arrangement_exert_policy_grants_total", help: "Exertion policy grants, by reason.", const_labels: {"reason" => "consolidation"}), + || EXERT_POLICY_CONSOLIDATION_GRANTS.load(Ordering::Relaxed), + ); +} + mod alloc { /// A Timely communication refill function that uses lgalloc. /// diff --git a/src/compute-types/src/dyncfgs.rs b/src/compute-types/src/dyncfgs.rs index 01cf2f5292d9c..c6173d99c53ae 100644 --- a/src/compute-types/src/dyncfgs.rs +++ b/src/compute-types/src/dyncfgs.rs @@ -134,6 +134,15 @@ pub const COLUMN_CHUNK_COMPRESS_MIN_DEPTH: Config = Config::new( ParameterScope::Replica, ); +/// Bypass resident pool slots when committing compression-eligible chunk bodies. +pub const ENABLE_COLUMN_CHUNK_DIRECT_COMPRESSED_OUTPUT: Config = Config::new( + "enable_column_chunk_direct_compressed_output", + false, + "Compress spilled chunk output directly into extents without resident slot admission. \ + Applies at the compression depth floor and above, at each spill.", + ParameterScope::Replica, +); + /// Resident-bytes budget fraction for chunk spilling. Two consumers read /// it: the column pager's tiered policy multiplies it against the /// announced memory limit, and the buffer pool (`mz_ore::pool`) @@ -223,6 +232,18 @@ pub const COLUMN_PAGED_BATCHER_EAGER_BACKING: Config = Config::new( ParameterScope::Replica, ); +/// Serve the buffer pool's nonresident reads on the async runtime's blocking +/// executor instead of the spill threads. Off, reads queue ahead of evictions +/// on the spill threads and need no runtime. Only meaningful with spill +/// workers; without them reads always use the runtime. +pub const COLUMN_POOL_RUNTIME_READS: Config = Config::new( + "column_pool_runtime_reads", + false, + "Serve buffer-pool nonresident reads on the async runtime's blocking executor rather than \ + the spill threads. Only meaningful with spill workers.", + ParameterScope::Replica, +); + /// Ceiling on the buffer pool's total RSS, as a fraction of *physical RAM* /// (never the announced limit, which includes swap on swap-provisioned /// nodes). The compressed-but-resident extent tier is the headroom above the @@ -850,6 +871,8 @@ pub fn all_dyncfgs(configs: ConfigSet) -> ConfigSet { .add(&COLUMN_PAGED_BATCHER_SWAP_PAGEOUT) .add(&COLUMN_PAGED_BATCHER_SPILL_WORKER_COUNT) .add(&COLUMN_PAGED_BATCHER_EAGER_BACKING) + .add(&COLUMN_POOL_RUNTIME_READS) .add(&COLUMN_PAGED_BATCHER_POOL_RSS_TARGET_FRACTION) .add(&COLUMN_CHUNK_COMPRESS_MIN_DEPTH) + .add(&ENABLE_COLUMN_CHUNK_DIRECT_COMPRESSED_OUTPUT) } diff --git a/src/compute/src/compute_state.rs b/src/compute/src/compute_state.rs index b19c8d8044c65..1a309c8c2b1b6 100644 --- a/src/compute/src/compute_state.rs +++ b/src/compute/src/compute_state.rs @@ -532,6 +532,7 @@ impl ComputeState { } else { let spill_threads = COLUMN_PAGED_BATCHER_SPILL_WORKER_COUNT.get(config); let eager_backing = COLUMN_PAGED_BATCHER_EAGER_BACKING.get(config); + let runtime_reads = COLUMN_POOL_RUNTIME_READS.get(config); // Budget derivation: fraction of physical RAM, with a 128 MiB // floor so the no-pressure case doesn't page per chunk. @@ -556,6 +557,7 @@ impl ComputeState { budget_bytes: total, spill_threads, eager_backing, + runtime_reads, rss_target_bytes: rss_target, }); if applied { @@ -567,6 +569,7 @@ impl ComputeState { budget_bytes = total, spill_threads, eager_backing, + runtime_reads, rss_target_bytes = rss_target, "chunk spill: applying buffer-pool config", ); @@ -581,6 +584,9 @@ impl ComputeState { let compress_min_depth = u8::try_from(COLUMN_CHUNK_COMPRESS_MIN_DEPTH.get(config)).unwrap_or(u8::MAX); mz_timely_util::columnar::chunk::set_compress_min_depth(compress_min_depth); + mz_timely_util::columnar::chunk::set_direct_compressed_output( + ENABLE_COLUMN_CHUNK_DIRECT_COMPRESSED_OUTPUT.get(config), + ); } // Remember the maintenance interval locally to avoid reading it from the config set on diff --git a/src/compute/src/server.rs b/src/compute/src/server.rs index b062f58da1adc..ced444311f74c 100644 --- a/src/compute/src/server.rs +++ b/src/compute/src/server.rs @@ -175,6 +175,7 @@ pub async fn serve( mz_timely_util::column_pager::tiered_policy(), ); mz_timely_util::pool_config::metrics::register(metrics_registry); + mz_cluster::client::register_exert_policy_metrics(metrics_registry); let config = Config { persist_clients, diff --git a/src/controller-types/src/dyncfgs.rs b/src/controller-types/src/dyncfgs.rs index 7a37b862df65f..774273d716770 100644 --- a/src/controller-types/src/dyncfgs.rs +++ b/src/controller-types/src/dyncfgs.rs @@ -76,6 +76,17 @@ pub const ARRANGEMENT_EXERT_PROPORTIONALITY: Config = Config::new( ParameterScope::Replica, ); +/// The proportionality halves per layer below the largest, so 16 lets the +/// policy request consolidation of layers within four levels of the largest +/// batch, the same reach compute uses. Larger values let every small published +/// batch be lifted into the largest batch on idle turns. +pub const STORAGE_ARRANGEMENT_EXERT_PROPORTIONALITY: Config = Config::new( + "storage_arrangement_exert_proportionality", + 16, + "Storage arrangement maintenance proportionality. Zero disables optional maintenance. Applies when a replica is provisioned.", + ParameterScope::Replica, +); + pub const ENABLE_PAUSED_CLUSTER_READHOLD_DOWNGRADE: Config = Config::new( "enable_paused_cluster_readhold_downgrade", true, @@ -94,5 +105,6 @@ pub fn all_dyncfgs(configs: ConfigSet) -> ConfigSet { .add(&ENABLE_TIMELY_ZERO_COPY_LGALLOC) .add(&TIMELY_ZERO_COPY_LIMIT) .add(&ARRANGEMENT_EXERT_PROPORTIONALITY) + .add(&STORAGE_ARRANGEMENT_EXERT_PROPORTIONALITY) .add(&ENABLE_PAUSED_CLUSTER_READHOLD_DOWNGRADE) } diff --git a/src/controller/src/clusters.rs b/src/controller/src/clusters.rs index 3af30e43670f6..6482f7e2bb70d 100644 --- a/src/controller/src/clusters.rs +++ b/src/controller/src/clusters.rs @@ -26,7 +26,8 @@ use mz_compute_client::logging::LogVariant; use mz_compute_types::config::{ComputeReplicaConfig, ComputeReplicaLogging}; use mz_controller_types::dyncfgs::{ ARRANGEMENT_EXERT_PROPORTIONALITY, CONTROLLER_PAST_GENERATION_REPLICA_CLEANUP_RETRY_INTERVAL, - ENABLE_TIMELY_ZERO_COPY, ENABLE_TIMELY_ZERO_COPY_LGALLOC, TIMELY_ZERO_COPY_LIMIT, + ENABLE_TIMELY_ZERO_COPY, ENABLE_TIMELY_ZERO_COPY_LGALLOC, + STORAGE_ARRANGEMENT_EXERT_PROPORTIONALITY, TIMELY_ZERO_COPY_LIMIT, }; use mz_controller_types::{ClusterId, ReplicaId}; use mz_orchestrator::NamespacedOrchestrator; @@ -704,11 +705,6 @@ impl Controller { let persist_pubsub_url = self.persist_pubsub_url.clone(); let secrets_args = self.secrets_args.to_flags(); - // TODO(teskje): use the same values as for compute? - let storage_proto_timely_config = TimelyConfig { - arrangement_exert_proportionality: 1337, - ..Default::default() - }; // These configure the replica's process rather than environmentd's, so // they are `ParameterScope::Replica` and must be read through this // replica's scoped overrides. They are baked into the process @@ -716,6 +712,11 @@ impl Controller { // environment-wide value or the override reaches the replica only when // it is next provisioned. let overrides = self.replica_dyncfg_overrides.get(&replica_id); + let storage_proto_timely_config = TimelyConfig { + arrangement_exert_proportionality: STORAGE_ARRANGEMENT_EXERT_PROPORTIONALITY + .get_with_overrides(&self.dyncfg, overrides), + ..Default::default() + }; let compute_proto_timely_config = TimelyConfig { arrangement_exert_proportionality: ARRANGEMENT_EXERT_PROPORTIONALITY .get_with_overrides(&self.dyncfg, overrides), diff --git a/src/ore/src/pool.rs b/src/ore/src/pool.rs index e985f57b52029..20d98c03c2cd9 100644 --- a/src/ore/src/pool.rs +++ b/src/ore/src/pool.rs @@ -80,6 +80,11 @@ use crate::pool::region::{Region, SIZE_CLASSES}; /// NOTE: Seen OoMs with Miri since it actually allocates the capacity. const CLASS_CAPACITY_BYTES: usize = if cfg!(miri) { 16 << 20 } else { 1 << 40 }; +// At most eight queued or running reads per pool. Completed buffers belong +// to callers, whose staging limits must bound their retained memory. +#[cfg(feature = "async")] +const ASYNC_READ_CONCURRENCY: usize = 8; + /// A chunk-provided transform between a chunk's body bytes and the stored /// bytes its extent holds. The pool owns scheduling: spill threads, the /// residency state machine, cancellation, and the ledger. It invokes the @@ -194,6 +199,17 @@ enum Residency { pub struct PoolStats { /// Chunks inserted. pub inserts: u64, + /// Inserts written directly to an extent because resident admission was full. + pub direct_extent_inserts: u64, + /// Chunks inserted directly into an extent by caller request. + pub cold_inserts: u64, + /// Nonresident reads submitted off-worker, to the spill threads or the + /// blocking executor. + pub async_reads: u64, + /// Of `async_reads`, those served on the spill threads. + pub spill_reads: u64, + /// Submitted reads not yet completed. + pub async_reads_in_flight: u64, /// Chunks freed (handle dropped). pub frees: u64, /// Backing writes elided: chunks dead before their compression @@ -279,6 +295,11 @@ pub struct PoolStats { #[derive(Debug, Default)] struct Counters { + /// Evicted-chunk reads served by the spill threads rather than a runtime. + spill_reads: AtomicU64, + direct_extent_inserts: AtomicU64, + cold_inserts: AtomicU64, + async_reads: AtomicU64, inserts: AtomicU64, spill_scheduled: AtomicU64, spill_cancelled: AtomicU64, @@ -378,7 +399,13 @@ struct PoolInner { /// counter read still has its bytes enforced rather than dropped. enforce_pending: std::sync::atomic::AtomicBool, counters: Counters, + #[cfg(feature = "async")] + read_slots: Arc, spill: Spill, + /// Benchmark hook: microseconds each evicted-chunk read sleeps before + /// decoding, standing in for device latency the fixture cannot produce. + #[cfg(feature = "test")] + read_delay_micros: AtomicU64, } /// Hand-off point between budget enforcement and spill threads. Eviction I/O @@ -405,6 +432,16 @@ struct Spill { threads: AtomicU64, /// Queued plus currently-processing entries; `quiesce` waits on zero. in_flight: AtomicU64, + /// Evicted-chunk reads awaiting a spill thread. Pushed under the `queue` + /// lock so a parking worker's emptiness check and the notify cannot miss + /// each other. Served before evictions: a read unblocks a worker, an + /// eviction only frees memory that the budget already accounted for. + reads: Mutex>>, + /// Queued plus currently-copying reads. + reads_in_flight: AtomicU64, + /// Sends nonresident reads to the async runtime's blocking executor even + /// when spill threads exist; see [`Pool::set_runtime_reads`]. + runtime_reads: std::sync::atomic::AtomicBool, /// Test-only lifecycle: production spill threads are immortal (the pool /// is a process singleton), but Miri rejects a test binary exiting with /// live threads, so tests stop and join them. @@ -419,6 +456,103 @@ struct Spill { /// unbounded queue of still-resident chunks. const SPILL_IN_FLIGHT_MAX: usize = 64; +/// An evicted-chunk read handed to the spill threads. +/// +/// The job owns the handle until the copy completes, so a caller that drops +/// its [`QueuedRead`] cannot free the chunk under the reading thread. +struct ReadJob { + handle: Arc, + state: Mutex, +} + +#[derive(Default)] +struct ReadJobState { + result: Option>, + waker: Option, +} + +impl std::fmt::Debug for ReadJob { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("ReadJob").finish_non_exhaustive() + } +} + +/// Resolves with the chunk's words once a spill thread has copied them out. +/// +/// Needs no runtime: completion wakes whatever `Waker` last polled it. +#[derive(Debug)] +pub struct QueuedRead(Arc); + +impl std::future::Future for QueuedRead { + type Output = Vec; + + fn poll( + self: std::pin::Pin<&mut Self>, + cx: &mut std::task::Context<'_>, + ) -> std::task::Poll> { + let mut state = self.0.state.lock().expect("read job poisoned"); + if let Some(words) = state.result.take() { + return std::task::Poll::Ready(words); + } + state.waker = Some(cx.waker().clone()); + std::task::Poll::Pending + } +} + +impl PoolInner { + /// Whether evicted-chunk reads run on the spill threads. + fn spill_reads_enabled(&self) -> bool { + self.spill.enabled.load(Ordering::Relaxed) + && self.spill.threads.load(Ordering::Relaxed) > 0 + && !self.spill.runtime_reads.load(Ordering::Relaxed) + } + + fn queue_read(&self, handle: Arc) -> QueuedRead { + let job = Arc::new(ReadJob { + handle, + state: Mutex::new(ReadJobState::default()), + }); + self.spill.reads_in_flight.fetch_add(1, Ordering::Relaxed); + { + let _queue = self.spill_queue(); + self.spill + .reads + .lock() + .expect("read queue poisoned") + .push_back(Arc::clone(&job)); + } + self.spill.cv.notify_one(); + QueuedRead(job) + } + + /// Serves one queued read on the calling spill thread. + fn read_step(&self) -> bool { + let job = self + .spill + .reads + .lock() + .expect("read queue poisoned") + .pop_front(); + let Some(job) = job else { + return false; + }; + let mut words = Vec::new(); + job.handle.read_into(&mut words); + let waker = { + let mut state = job.state.lock().expect("read job poisoned"); + state.result = Some(words); + state.waker.take() + }; + // Decrement before waking so a caller that observes its result also + // observes no read in flight. + self.spill.reads_in_flight.fetch_sub(1, Ordering::Relaxed); + if let Some(waker) = waker { + waker.wake(); + } + true + } +} + /// What a spill thread does with a chunk once compressed. #[derive(Clone, Copy, PartialEq, Eq)] enum SpillKind { @@ -565,15 +699,31 @@ impl Pool { enforcing: Mutex::new(()), enforce_pending: std::sync::atomic::AtomicBool::new(false), counters: Counters::default(), + #[cfg(feature = "async")] + read_slots: Arc::new(tokio::sync::Semaphore::new(ASYNC_READ_CONCURRENCY)), spill: Spill::default(), + #[cfg(feature = "test")] + read_delay_micros: AtomicU64::new(0), }))) } + /// Benchmark hook: every evicted-chunk read sleeps `delay` before decoding. + #[cfg(feature = "test")] + pub fn set_read_delay(&self, delay: std::time::Duration) { + self.0.read_delay_micros.store( + u64::try_from(delay.as_micros()).unwrap_or(u64::MAX), + Ordering::Relaxed, + ); + } + /// Allocates a chunk of `len` words and fills it in place: `fill` - /// receives the chunk's slot memory directly and must overwrite all of - /// it (the slot's prior contents are unspecified), so serialization - /// writes its single copy straight into pool memory. The returned handle - /// starts `UnbackedResident`. A zero `len` returns a length-0 handle + /// receives `len` contiguous words and must overwrite all of them. + /// With resident admission, these are the slot's unspecified prior + /// contents, so serialization writes directly into pool memory. + /// Otherwise, `fill` writes into staging for a synchronous extent write. + /// The returned handle + /// starts `UnbackedResident` when admission has room, otherwise `Evicted` + /// with a directly written extent. A zero `len` returns a length-0 handle /// holding no slot; payloads beyond the largest size class fall back to /// a plain heap allocation, always resident, a prototype limitation. /// `hints` steer eviction and write-behind policy; callers without @@ -620,14 +770,22 @@ impl Pool { .oversize_payloads .fetch_add(1, Ordering::Relaxed); } + if class.is_some() { + if !inner.reserve_insert(len_bytes) { + inner.enforce_budget(); + if !inner.reserve_insert(len_bytes) { + return self.insert_extent(len, hints, codec, fill); + } + } + } else { + inner + .counters + .resident_bytes + .fetch_add(u64::cast_from(len_bytes), Ordering::Relaxed); + } // A class with no free slot degrades to the heap path below: an // unpageable chunk beats a dead replica. let slot = class.and_then(|class| inner.alloc_slot(class, len_bytes)); - // Whichever home the payload found, it is resident. - inner - .counters - .resident_bytes - .fetch_add(u64::cast_from(len_bytes), Ordering::Relaxed); let meta = match (class, slot) { (Some(class), Some(slot)) => { let region = &inner.regions[class]; @@ -689,11 +847,98 @@ impl Pool { ChunkHandle { meta } } + /// Encode `data` directly into an extent without allocating a resident slot. + /// + /// Compression is synchronous. The extent participates in the pool's + /// compressed-residency accounting and reclamation. Empty and oversize + /// payloads use the same fallback as [`Pool::insert_with`]. + pub fn insert_cold( + &self, + data: &[u64], + hints: ChunkHints, + codec: &'static dyn ExtentCodec, + ) -> ChunkHandle { + if data.is_empty() || region::size_class_for(std::mem::size_of_val(data)).is_none() { + return self.insert_with(data.len(), hints, codec, |dst| dst.copy_from_slice(data)); + } + let inner = &self.0; + let extent = SwapExtent::write(&inner.extent_arena, data, codec, Scratch::Shrink); + inner.counters.inserts.fetch_add(1, Ordering::Relaxed); + inner.counters.cold_inserts.fetch_add(1, Ordering::Relaxed); + self.finish_extent(data.len(), hints, codec, extent) + } + + fn insert_extent( + &self, + len: usize, + hints: ChunkHints, + codec: &'static dyn ExtentCodec, + fill: impl FnOnce(&mut [u64]), + ) -> ChunkHandle { + // Compress synchronously to keep denied insertions from queuing + // uncompressed payloads behind an occupied enforcer. + let mut words = vec![0; len]; + fill(&mut words); + let inner = &self.0; + let extent = SwapExtent::write(&inner.extent_arena, &words, codec, Scratch::Shrink); + drop(words); + inner + .counters + .direct_extent_inserts + .fetch_add(1, Ordering::Relaxed); + self.finish_extent(len, hints, codec, extent) + } + + fn finish_extent( + &self, + len: usize, + hints: ChunkHints, + codec: &'static dyn ExtentCodec, + extent: SwapExtent, + ) -> ChunkHandle { + let inner = &self.0; + let meta = Arc::new(ChunkMeta::new( + inner, + len, + region::size_class_for(len * 8), + hints.depth, + codec, + Residency::Evicted, + None, + None, + )); + inner.live_chunks.fetch_add(1, Ordering::Relaxed); + { + let mut state = meta.state(); + inner.commit_extent(&meta, &mut state, extent); + } + inner.enforce_or_defer_compressed_cap(); + ChunkHandle { meta } + } + /// Snapshot of the pool's counters. pub fn stats(&self) -> PoolStats { let c = &self.0.counters; PoolStats { inserts: c.inserts.load(Ordering::Relaxed), + direct_extent_inserts: c.direct_extent_inserts.load(Ordering::Relaxed), + cold_inserts: c.cold_inserts.load(Ordering::Relaxed), + async_reads: c.async_reads.load(Ordering::Relaxed), + spill_reads: c.spill_reads.load(Ordering::Relaxed), + async_reads_in_flight: { + let queued = self.0.spill.reads_in_flight.load(Ordering::Relaxed); + #[cfg(feature = "async")] + { + queued + + u64::cast_from( + ASYNC_READ_CONCURRENCY - self.0.read_slots.available_permits(), + ) + } + #[cfg(not(feature = "async"))] + { + queued + } + }, frees: c.frees.load(Ordering::Relaxed), writes_elided: c.writes_elided.load(Ordering::Relaxed), evictions_compress: c.evictions_compress.load(Ordering::Relaxed), @@ -762,6 +1007,14 @@ impl Pool { self.0.spill.enabled.store(true, Ordering::Relaxed); } + /// Routes nonresident reads to the async runtime's blocking executor + /// instead of the spill threads. Off by default: with spill threads + /// spawned, reads queue ahead of evictions on those threads and need no + /// runtime. Only meaningful with spill threads spawned. + pub fn set_runtime_reads(&self, runtime: bool) { + self.0.spill.runtime_reads.store(runtime, Ordering::Relaxed); + } + /// Enables or disables eager backing: when on, idle spill threads /// compress unbacked chunks to `BackedResident` ahead of pressure, so /// budget-driven eviction becomes a pure page release. Costs CPU on @@ -790,6 +1043,15 @@ impl Pool { } } + /// Test hook: waits until no queued read remains, so tests observe the + /// job of a dropped [`QueuedRead`] finishing. + #[cfg(test)] + fn quiesce_reads(&self) { + while self.0.spill.reads_in_flight.load(Ordering::Relaxed) > 0 { + std::thread::yield_now(); + } + } + /// Test hook: stops and joins the spill threads, so a test binary exits /// with none alive (which Miri requires). Stopped threads process no /// further queued work; call [`Pool::quiesce_spill`] first when the test @@ -986,6 +1248,23 @@ impl PoolInner { } } + /// Reserve insertion bytes before populating a slot. Enforcement can + /// lag by one eighth of the budget, or one payload for small budgets. + /// Read admissions use the budget itself and cannot consume this slack. + fn reserve_insert(&self, len_bytes: usize) -> bool { + let len = u64::cast_from(len_bytes); + self.counters + .resident_bytes + .try_update(Ordering::Relaxed, Ordering::Relaxed, |cur| { + let budget = self.budget_bytes.load(Ordering::Relaxed); + let ceiling = budget.saturating_add((budget / 8).max(len)); + let next = cur.checked_add(len)?; + let oversize = self.counters.oversize_bytes.load(Ordering::Relaxed); + (next.saturating_sub(oversize) <= ceiling).then_some(next) + }) + .is_ok() + } + fn enforce_budget(&self) { // Single-flight: enforcement runs synchronously on whichever thread // trips it (every insert), and concurrent passes would @@ -1178,6 +1457,9 @@ impl PoolInner { // loop (job completion, condvar wakeup, park timeout) trims the // compressed tier if needed. A single atomic load when under cap. self.enforce_compressed_cap(); + if self.read_step() { + continue; + } let popped = self.spill_queue().pop_front(); if let Some(meta) = popped { self.spill_process(&meta, SpillKind::Evict); @@ -1187,12 +1469,19 @@ impl PoolInner { if self.spill.eager.load(Ordering::Relaxed) && self.back_one() { continue; } - // Nothing to evict or back: park. Re-checking emptiness under - // the queue lock closes the lost-wakeup window (hand-offs push - // under this lock before notifying); the timeout backstops - // everything else (fresh inserts, tier growth, lost notifies). + // Nothing to read, evict, or back: park. Re-checking emptiness + // under the queue lock closes the lost-wakeup window (hand-offs + // and reads push under this lock before notifying); the timeout + // backstops everything else (fresh inserts, tier growth, lost + // notifies). let queue = self.spill_queue(); - if queue.is_empty() { + let reads_empty = self + .spill + .reads + .lock() + .expect("read queue poisoned") + .is_empty(); + if queue.is_empty() && reads_empty { let _ = self .spill .cv @@ -1889,6 +2178,70 @@ impl ChunkHandle { self.read_impl(0..self.meta.len, dst, false); } + /// Copy this chunk without admitting it to the pool, offloading nonresident reads. + /// + /// Resident slots are copied inline if their state lock is available. Other + /// reads require a Tokio runtime. A submitted read retains its handle and + /// concurrency permit until it finishes, even if the caller cancels. + /// The returned buffer belongs to the caller. + #[cfg(feature = "async")] + pub async fn read_async(self: &Arc) -> Vec { + if let Some(words) = self.try_read_resident() { + return words; + } + let pool = &self.meta.pool; + if pool.spill_reads_enabled() { + pool.counters.async_reads.fetch_add(1, Ordering::Relaxed); + pool.counters.spill_reads.fetch_add(1, Ordering::Relaxed); + return pool.queue_read(Arc::clone(self)).await; + } + let permit = Arc::clone(&self.meta.pool.read_slots) + .acquire_owned() + .await + .expect("pool read semaphore remains open"); + let handle = Arc::clone(self); + self.meta + .pool + .counters + .async_reads + .fetch_add(1, Ordering::Relaxed); + crate::task::spawn_blocking( + || "pool_read", + move || { + // A cancelled JoinHandle must not release admission while its + // blocking read still owns a slot or extent reference. + let _permit = permit; + let mut words = Vec::new(); + handle.read_into(&mut words); + words + }, + ) + .await + } + + #[cfg(feature = "async")] + fn try_read_resident(&self) -> Option> { + if self.meta.len == 0 { + return Some(Vec::new()); + } + let mut state = match self.meta.state.try_lock() { + Ok(state) => state, + Err(std::sync::TryLockError::WouldBlock) => return None, + Err(std::sync::TryLockError::Poisoned(_)) => panic!("chunk state poisoned"), + }; + match state.residency { + Residency::UnbackedResident | Residency::BackedResident | Residency::WriteInFlight => { + state.touched = true; + let slot = state.slot.expect("resident non-empty chunk has a slot"); + // SAFETY: the state lock prevents eviction from releasing this slot. + let words = unsafe { self.meta.pool.slot_data(&self.meta, slot) }; + Some(words.to_vec()) + } + // Oversize heap copies have no slot-size bound, so keep them off-worker. + Residency::Oversize | Residency::Evicted => None, + } + } + /// As [`ChunkHandle::read_into`], restricted to the word range `range` /// of the chunk's contents, which must lie within them. `dst` receives /// exactly the range. @@ -1951,6 +2304,13 @@ impl ChunkHandle { dst.extend_from_slice(&payload[range]); } Residency::Evicted => { + #[cfg(feature = "test")] + { + let micros = meta.pool.read_delay_micros.load(Ordering::Relaxed); + if micros > 0 { + std::thread::sleep(std::time::Duration::from_micros(micros)); + } + } let slot = if admit { meta.pool.admit_slot(meta) } else { @@ -2147,6 +2507,13 @@ impl Drop for ChunkHandle { #[cfg(test)] mod tests { + #[cfg(feature = "async")] + use std::future::Future; + #[cfg(feature = "async")] + use std::sync::atomic::{AtomicBool, AtomicUsize}; + #[cfg(feature = "async")] + use std::task::{Context, Poll, Waker}; + use super::*; use crate::pool::extent::TEST_CODEC; @@ -2162,6 +2529,223 @@ mod tests { pool } + #[cfg(feature = "async")] + #[derive(Debug, Default)] + struct DelayedReadCodec { + entered: AtomicUsize, + released: AtomicBool, + } + + #[cfg(feature = "async")] + impl ExtentCodec for DelayedReadCodec { + fn encode(&self, body: &[u8], out: &mut Vec) { + TEST_CODEC.encode(body, out); + } + + fn decode(&self, stored: &[u8], body: &mut [u8]) { + // NOTE: Test-only delay models a page fault while the state lock is held. + self.entered.fetch_add(1, Ordering::SeqCst); + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(30); + while !self.released.load(Ordering::SeqCst) { + assert!( + std::time::Instant::now() < deadline, + "test read was released" + ); + std::thread::sleep(std::time::Duration::from_millis(1)); + } + TEST_CODEC.decode(stored, body); + } + } + + #[cfg(feature = "async")] + struct ReleaseReadsOnDrop(&'static DelayedReadCodec); + + #[cfg(feature = "async")] + impl Drop for ReleaseReadsOnDrop { + fn drop(&mut self) { + self.0.released.store(true, Ordering::SeqCst); + } + } + + #[cfg(feature = "async")] + async fn wait_for_read_state(mut ready: impl FnMut() -> bool) { + tokio::time::timeout(std::time::Duration::from_secs(10), async { + while !ready() { + tokio::task::yield_now().await; + } + }) + .await + .expect("read workers made progress"); + } + + #[cfg(feature = "async")] + #[mz_ore::test(tokio::test)] + #[cfg_attr(miri, ignore)] + async fn async_read_round_trip_without_admission() { + for budget in [0, usize::MAX] { + let pool = test_pool(budget); + let expected = payload(8192, 42); + let handle = Arc::new(insert(&pool, &mut expected.clone())); + let resident = pool.stats().resident_bytes; + assert_eq!(handle.read_async().await, expected); + assert_eq!(pool.stats().resident_bytes, resident); + assert_eq!(pool.stats().async_reads, u64::from(budget == 0)); + assert_eq!(pool.stats().async_reads_in_flight, 0); + } + } + + #[cfg(feature = "async")] + #[mz_ore::test] + fn resident_async_reads_complete_without_a_runtime() { + let pool = test_pool(usize::MAX); + let expected = payload(8192, 42); + let handle = Arc::new(insert(&pool, &mut expected.clone())); + let mut read = Box::pin(handle.read_async()); + assert_eq!( + read.as_mut().poll(&mut Context::from_waker(Waker::noop())), + Poll::Ready(expected) + ); + assert_eq!(pool.stats().async_reads, 0); + pool.evict(&handle); + assert!(handle.try_read_resident().is_none()); + } + + #[cfg(feature = "async")] + #[mz_ore::test] + #[cfg_attr(miri, ignore)] + fn evicted_reads_complete_on_spill_threads_without_a_runtime() { + struct ChannelWake(std::sync::mpsc::Sender<()>); + impl std::task::Wake for ChannelWake { + fn wake(self: Arc) { + let _ = self.0.send(()); + } + } + let pool = test_pool(0); + pool.set_spill_threads(2); + let expected = payload(8192, 42); + let handle = Arc::new(insert(&pool, &mut expected.clone())); + pool.quiesce_spill(); + assert!(handle.try_read_resident().is_none(), "a zero budget evicts"); + let (sender, receiver) = std::sync::mpsc::channel(); + let waker = Waker::from(Arc::new(ChannelWake(sender))); + let mut cx = Context::from_waker(&waker); + let mut queued = Box::pin(handle.read_async()); + let words = loop { + match queued.as_mut().poll(&mut cx) { + Poll::Ready(words) => break words, + Poll::Pending => receiver + .recv_timeout(std::time::Duration::from_secs(10)) + .expect("spill thread wakes the reader"), + } + }; + assert_eq!(words, expected); + let stats = pool.stats(); + assert_eq!(stats.spill_reads, 1); + assert_eq!(stats.async_reads, 1); + assert_eq!(stats.async_reads_in_flight, 0); + // A dropped future leaves its job to finish on the spill thread. + let mut dropped = Box::pin(handle.read_async()); + let _ = dropped.as_mut().poll(&mut cx); + drop(dropped); + pool.quiesce_reads(); + assert_eq!(pool.stats().spill_reads, 2); + assert_eq!(pool.stats().async_reads_in_flight, 0); + assert_eq!( + read(&handle), + expected, + "the chunk outlived the dropped read" + ); + pool.join_spill_threads(); + } + + #[cfg(feature = "async")] + #[mz_ore::test(tokio::test)] + #[cfg_attr(miri, ignore)] + async fn contended_resident_reads_yield() { + let pool = test_pool(usize::MAX); + let expected = payload(8192, 42); + let handle = Arc::new(insert(&pool, &mut expected.clone())); + let mut read = Box::pin(handle.read_async()); + { + let _state = handle.meta.state(); + assert!( + read.as_mut() + .poll(&mut Context::from_waker(Waker::noop())) + .is_pending() + ); + } + assert_eq!(read.await, expected); + assert_eq!(pool.stats().async_reads, 1); + } + + #[cfg(feature = "async")] + #[mz_ore::test(tokio::test)] + #[ignore = "local resident-read dispatch comparison"] + async fn resident_read_microbench() { + for words in [8 << 10, 256 << 10] { + let pool = test_pool(usize::MAX); + let expected = payload(words, 42); + let handle = Arc::new(insert(&pool, &mut expected.clone())); + assert_eq!(handle.read_async().await, expected); + let start = std::time::Instant::now(); + for _ in 0..1024 { + std::hint::black_box(handle.read_async().await); + } + eprintln!( + "RESIDENT_READ bytes={} iterations=1024 us={} offloaded={}", + words * 8, + start.elapsed().as_micros(), + pool.stats().async_reads, + ); + } + } + + #[cfg(feature = "async")] + #[mz_ore::test(tokio::test)] + #[cfg_attr(miri, ignore)] + async fn async_reads_remain_bounded_after_cancellation() { + let pool = test_pool(0); + let codec = Box::leak(Box::new(DelayedReadCodec::default())); + let release = ReleaseReadsOnDrop(codec); + let expected = payload(8192, 7); + let handles: Vec<_> = (0..=ASYNC_READ_CONCURRENCY) + .map(|_| { + Arc::new( + pool.insert_with(expected.len(), ChunkHints::default(), codec, |dst| { + dst.copy_from_slice(&expected); + }), + ) + }) + .collect(); + let mut reads: Vec<_> = handles.iter().map(|h| Box::pin(h.read_async())).collect(); + for read in &mut reads { + assert!(futures::poll!(read.as_mut()).is_pending()); + } + wait_for_read_state(|| codec.entered.load(Ordering::SeqCst) == ASYNC_READ_CONCURRENCY) + .await; + assert_eq!( + pool.stats().async_reads, + u64::cast_from(ASYNC_READ_CONCURRENCY) + ); + assert_eq!( + pool.stats().async_reads_in_flight, + u64::cast_from(ASYNC_READ_CONCURRENCY) + ); + drop(reads); + drop(handles); + // The unsubmitted ninth read frees immediately. Detached blocking + // jobs must retain both their handles and their admission permits. + assert_eq!(pool.stats().frees, 1); + assert_eq!( + pool.stats().async_reads_in_flight, + u64::cast_from(ASYNC_READ_CONCURRENCY) + ); + drop(release); + wait_for_read_state(|| pool.stats().frees == u64::cast_from(ASYNC_READ_CONCURRENCY + 1)) + .await; + assert_eq!(pool.stats().async_reads_in_flight, 0); + } + /// Scales an iteration count down under Miri, where one interpreted /// compression costs what thousands do natively. fn rounds(native: u64, miri: u64) -> u64 { @@ -3172,6 +3756,67 @@ mod tests { ); } + #[mz_ore::test] + fn insertion_debt_is_bounded_during_enforcement() { + let budget = 2 * SMALL * 8; + let pool = test_pool(budget); + let guard = pool.0.enforcing.lock().expect("enforcement lock"); + let mut handles = Vec::new(); + for seed in 0..16 { + handles.push(insert(&pool, &mut payload(SMALL, seed))); + assert!( + pool.stats().resident_bytes <= u64::cast_from(budget + SMALL * 8), + "an occupied enforcer must not allow unlimited insertion debt", + ); + } + drop(guard); + for (seed, handle) in handles.iter().enumerate() { + assert_eq!(read(handle), payload(SMALL, u64::cast_from(seed))); + } + drop(handles); + assert_eq!(pool.stats().resident_bytes, 0); + assert_eq!(pool.stats().live_chunks, 0); + assert_eq!(pool.stats().extent_resident_bytes, 0); + } + + #[mz_ore::test] + fn admission_reserves_before_concurrent_fills() { + let budget = 2 * SMALL * 8; + let pool = test_pool(budget); + let guard = pool.0.enforcing.lock().expect("enforcement lock"); + let gate = Arc::new(std::sync::Barrier::new(9)); + let threads: Vec<_> = (0..8u64) + .map(|seed| { + let pool = pool.clone(); + let gate = Arc::clone(&gate); + std::thread::spawn(move || { + pool.insert_with(SMALL, ChunkHints::default(), &TEST_CODEC, |dst| { + gate.wait(); + gate.wait(); + dst.copy_from_slice(&payload(SMALL, seed)); + }) + }) + }) + .collect(); + gate.wait(); + let reserved = pool.stats().resident_bytes; + // Release every producer even if the assertion fails. + gate.wait(); + let handles: Vec<_> = threads + .into_iter() + .map(|t| t.join().expect("producer panicked")) + .collect(); + drop(guard); + assert!(reserved <= u64::cast_from(budget + SMALL * 8)); + assert!(pool.stats().direct_extent_inserts > 0); + for (seed, handle) in handles.iter().enumerate() { + assert_eq!(read(handle), payload(SMALL, u64::cast_from(seed))); + } + drop(handles); + assert_eq!(pool.stats().resident_bytes, 0); + assert_eq!(pool.stats().live_chunks, 0); + } + #[mz_ore::test] fn set_budget_retunes_in_place() { let pool = test_pool(usize::MAX); @@ -3640,6 +4285,61 @@ mod tests { assert_eq!(range, want[8..24], "range reads copy the range directly"); } + #[mz_ore::test] + fn cold_insert_skips_resident_admission() { + let codecs: [&'static dyn ExtentCodec; 2] = [&TEST_CODEC, &IDENTITY_CODEC]; + for codec in codecs { + let pool = test_pool(256 << 20); + pool.set_rss_target(1 << 30); + let want = payload(SMALL, 703); + let handle = pool.insert_cold(&want, ChunkHints { depth: 3 }, codec); + assert_eq!(handle.residency(), Residency::Evicted); + assert_eq!(handle.meta.depth, 3); + let stats = pool.stats(); + assert_eq!(stats.inserts, 1); + assert_eq!(stats.cold_inserts, 1); + assert_eq!(stats.direct_extent_inserts, 0); + assert_eq!(stats.resident_bytes, 0); + assert_eq!(stats.live_chunks, 1); + assert!(stats.extent_resident_bytes > 0); + assert_eq!(read(&handle), want); + assert_eq!(pool.stats().resident_bytes, 0); + let mut range = Vec::new(); + handle.read_range_into(3..11, &mut range); + assert_eq!(range, want[3..11]); + assert_eq!(read_admit(&handle), want); + assert!(pool.stats().resident_bytes > 0); + pool.evict(&handle); + assert_eq!( + pool.stats().extent_bytes_written, + stats.extent_bytes_written + ); + assert_eq!(read(&handle), want); + drop(handle); + let stats = pool.stats(); + assert_eq!(stats.resident_bytes, 0); + assert_eq!(stats.live_chunks, 0); + assert_eq!(stats.extent_resident_bytes, 0); + assert_eq!(stats.frees, 1); + } + } + + #[mz_ore::test] + #[cfg_attr(miri, ignore)] + fn cold_insert_empty_and_oversize_fallbacks() { + let pool = test_pool(usize::MAX); + let empty = pool.insert_cold(&[], ChunkHints::default(), &TEST_CODEC); + assert_eq!(read(&empty), Vec::::new()); + let want = payload(SIZE_CLASSES[SIZE_CLASSES.len() - 1] / 8 + 1, 704); + let big = pool.insert_cold(&want, ChunkHints::default(), &TEST_CODEC); + assert_eq!(big.residency(), Residency::Oversize); + assert_eq!(read(&big), want); + assert_eq!(pool.stats().inserts, 2); + assert_eq!(pool.stats().cold_inserts, 0); + drop(big); + assert_eq!(pool.stats().resident_bytes, 0); + } + #[mz_ore::test] fn insert_with_fills_in_place() { let pool = test_pool(usize::MAX); diff --git a/src/storage-types/src/dyncfgs.rs b/src/storage-types/src/dyncfgs.rs index a98278efa5621..960bc2a89d669 100644 --- a/src/storage-types/src/dyncfgs.rs +++ b/src/storage-types/src/dyncfgs.rs @@ -442,6 +442,35 @@ pub const ENABLE_UPSERT_CHUNKED_STASH: Config = Config::new( ParameterScope::Replica, ); +/// Offload chunked upsert stash drains and feedback lookups to the blocking +/// executor. Merge reads remain synchronous. Read at operator construction. +pub const ENABLE_UPSERT_ASYNC_READS: Config = Config::new( + "enable_upsert_async_reads", + false, + "Read sealed stash and feedback chunks asynchronously when enable_upsert_v2 and \ + enable_upsert_chunked_stash are true. Takes effect on new dataflows.", + ParameterScope::Replica, +); + +/// Yield during source-stash merges and feedback compaction reads. +/// Applies to newly constructed chunked upsert-v2 dataflows without payload separation. +/// Independent of drain/probe read offload in `ENABLE_UPSERT_ASYNC_READS`. +pub const ENABLE_UPSERT_ASYNC_MERGES: Config = Config::new( + "enable_upsert_async_merges", + false, + "Use resumable source and feedback merges in chunked upsert-v2. Takes effect on new dataflows. Decoded merge inputs share a 256 MiB process budget.", + ParameterScope::Replica, +); + +/// Separate upsert-v2 payload blocks from columnar merge metadata. +/// Read once per dataflow. Uses the process pool and asynchronous payload reads. +pub const ENABLE_UPSERT_PAYLOAD_STASH: Config = Config::new( + "enable_upsert_payload_stash", + false, + "Use payload-separated columnar state for upsert-v2. Takes effect on new dataflows.", + ParameterScope::Replica, +); + // RocksDB /// How many times to try to cleanup old RocksDB DB's on disk before giving up. @@ -566,6 +595,9 @@ pub fn all_dyncfgs(configs: ConfigSet) -> ConfigSet { .add(&SUSPENDABLE_SOURCES) .add(&ENABLE_UPSERT_PAGED_SPILL) .add(&ENABLE_UPSERT_CHUNKED_STASH) + .add(&ENABLE_UPSERT_ASYNC_READS) + .add(&ENABLE_UPSERT_ASYNC_MERGES) + .add(&ENABLE_UPSERT_PAYLOAD_STASH) .add(&WALLCLOCK_GLOBAL_LAG_HISTOGRAM_RETENTION_INTERVAL) .add(&WALLCLOCK_LAG_HISTORY_RETENTION_INTERVAL) .add(&crate::sources::sql_server::CDC_CLEANUP_CHANGE_TABLE) diff --git a/src/storage/src/upsert_continual_feedback_v2.rs b/src/storage/src/upsert_continual_feedback_v2.rs index 5b4b699ac2048..8a23baf21c983 100644 --- a/src/storage/src/upsert_continual_feedback_v2.rs +++ b/src/storage/src/upsert_continual_feedback_v2.rs @@ -63,10 +63,13 @@ //! //! ## Stash flavors //! -//! [`UpsertStashFlavor`], resolved from `enable_upsert_chunked_stash` at -//! operator construction, selects between two instantiations of the same +//! [`UpsertStashFlavor`], resolved from the replica configuration at +//! operator construction, selects between three instantiations of the same //! loop: //! +//! * **Payload**: columnar chunks contain metadata and payload locators. +//! Separate manifests own pool-backed rows, and feedback uses exact, +//! operator-local payload identities for consolidation. //! * **Chunked**: the stash is differential's chunk merge batcher over //! `ColumnChunk`s and the feedback arrangement is a spine of chunk batches. //! Committed chunk bodies spill to the process buffer pool, the drain @@ -77,7 +80,9 @@ //! through the storage-owned column pager, and prior state comes back //! through a trace cursor. //! -//! Both flavors' spill paths are gated by `enable_upsert_paged_spill`. +//! Paged and chunked spill paths are gated by `enable_upsert_paged_spill`. +//! The experimental payload flavor is selected by `enable_upsert_payload_stash` +//! and always allocates its external payloads through the process pool. //! //! ## Eligibility condition (total order) //! @@ -88,6 +93,8 @@ //! with `p < ts` is ineligible (persist hasn't caught up), and one with //! `ts < p` is already persisted and dropped. +mod payload; + use std::fmt::Debug; use differential_dataflow::difference::{IsZero, Semigroup}; @@ -97,7 +104,7 @@ use differential_dataflow::logging::Logger; use differential_dataflow::operators::arrange::agent::TraceAgent; use differential_dataflow::operators::arrange::arrangement::{Arranged, arrange_core}; use differential_dataflow::trace::chunk::{ChunkBatcher, ChunkBuilder, ChunkSpine}; -use differential_dataflow::trace::{Batcher, Cursor, Description, TraceReader}; +use differential_dataflow::trace::{BatchReader, Batcher, Cursor, Description, TraceReader}; use differential_dataflow::{AsCollection, VecCollection}; use mz_dyncfg::ConfigSet; use mz_repr::{Datum, Diff, GlobalId, Row}; @@ -105,7 +112,10 @@ use mz_repr::{Datum, Diff, GlobalId, Row}; #[cfg(feature = "fuzzing")] use mz_row_spine::DatumSeq; use mz_row_spine::{ValRowColPagedBuilder, ValRowSpine}; -use mz_storage_types::dyncfgs::ENABLE_UPSERT_CHUNKED_STASH; +use mz_storage_types::dyncfgs::{ + ENABLE_UPSERT_ASYNC_MERGES, ENABLE_UPSERT_ASYNC_READS, ENABLE_UPSERT_CHUNKED_STASH, + ENABLE_UPSERT_PAYLOAD_STASH, +}; use mz_storage_types::errors::{DataflowError, EnvelopeError, UpsertError}; use mz_timely_util::builder_async::{ AsyncOutputHandle, Event as AsyncEvent, OperatorBuilder as AsyncOperatorBuilder, @@ -137,7 +147,7 @@ use crate::upsert::UpsertSourceTime; use crate::upsert::UpsertValue; /// Which stash and feedback-arrangement representation the upsert-v2 -/// operator instantiates. The two flavors run the same operator loop; they +/// operator instantiates. All flavors run the same operator loop. They /// differ in the batcher, the feedback trace, and how the drain reads prior /// state. See the module docs for the comparison. #[derive(Clone, Copy, Debug)] @@ -146,9 +156,13 @@ pub enum UpsertStashFlavor { /// arrangement, cursor-based drain. Spills through the storage-owned /// column pager. Paged, + /// Columnar metadata with independently owned payload blocks. + Payload, /// Chunk merge batcher stash, chunk-spine feedback arrangement, /// bulk-probe drain. Spills through the process buffer pool. - Chunked, + Chunked { async_reads: bool }, + /// Async source batching and feedback trace compaction. + Resumable { async_reads: bool }, } impl UpsertStashFlavor { @@ -156,8 +170,17 @@ impl UpsertStashFlavor { /// source at operator construction time, so a dataflow keeps one flavor /// for its whole life even if the flag flips underneath it. pub fn from_config(config: &ConfigSet) -> Self { - if ENABLE_UPSERT_CHUNKED_STASH.get(config) { - Self::Chunked + if ENABLE_UPSERT_PAYLOAD_STASH.get(config) { + Self::Payload + } else if ENABLE_UPSERT_CHUNKED_STASH.get(config) && ENABLE_UPSERT_ASYNC_MERGES.get(config) + { + Self::Resumable { + async_reads: ENABLE_UPSERT_ASYNC_READS.get(config), + } + } else if ENABLE_UPSERT_CHUNKED_STASH.get(config) { + Self::Chunked { + async_reads: ENABLE_UPSERT_ASYNC_READS.get(config), + } } else { Self::Paged } @@ -245,18 +268,18 @@ type FeedbackSpine = ChunkSpine>; // `(key, time)` runs. #[derive(Clone, Debug, Default, columnar::Columnar)] #[columnar(derive(PartialEq, Eq, PartialOrd, Ord))] -struct UpsertDiff { +struct UpsertDiff { from_time: O, - value: Option, + value: Option, } -impl IsZero for UpsertDiff { +impl IsZero for UpsertDiff { fn is_zero(&self) -> bool { false } } -impl Semigroup for UpsertDiff { +impl Semigroup for UpsertDiff { fn plus_equals(&mut self, rhs: &Self) { if rhs.from_time > self.from_time { *self = rhs.clone(); @@ -270,15 +293,16 @@ impl Semigroup for UpsertDiff { // wins" comparison — copying the value `Row` out of the column solely when `rhs` // wins. Losing folds (the common case for a repeatedly-updated key) then pay no // `Row` copy at all. -impl<'a, O> Semigroup>> for UpsertDiff +impl<'a, O, V> Semigroup>> for UpsertDiff where O: columnar::Columnar + Ord + Clone, + V: columnar::Columnar + Clone, { - fn plus_equals(&mut self, rhs: &columnar::Ref<'a, UpsertDiff>) { + fn plus_equals(&mut self, rhs: &columnar::Ref<'a, UpsertDiff>) { let rhs_from_time = ::into_owned(rhs.from_time); if rhs_from_time > self.from_time { self.from_time = rhs_from_time; - self.value = as columnar::Columnar>::into_owned(rhs.value); + self.value = as columnar::Columnar>::into_owned(rhs.value); } } } @@ -286,11 +310,11 @@ where /// One source-stash update: a key, its dataflow time, and the payload diff. /// `O` is the columnar order key projected from the source `FromTime` (see /// [`UpsertSourceTime`]). -type UpsertUpdate = (UpsertKey, T, UpsertDiff); +type UpsertUpdate = (UpsertKey, T, UpsertDiff); /// One stash chunk: a sorted, consolidated run of updates, resident or /// spilled to the buffer pool. -type UpsertChunk = ColumnChunk>; +type UpsertChunk = ColumnChunk>; /// The chunked flavor's stash: differential's chunk merge batcher over /// `ColumnChunk`s. Data is pushed in unsorted. The batcher maintains @@ -305,11 +329,11 @@ type UpsertChunkBatcher = ChunkBatcher>; /// like [`UpsertChunkBatcher`] but storing each chain entry as a `Column` /// routed through the storage-owned pager, which pages cold chains out of /// RSS. -type UpsertPagedBatcher = ColumnMergeBatcher>; +type UpsertPagedBatcher = ColumnMergeBatcher>; /// The chunker that sorts and consolidates raw input into the `Column` chunks /// both stash batchers consume. -type UpsertChunker = ColumnChunker>; +type UpsertChunker = ColumnChunker>; /// The operator's data-output handle. A fueled `Vec` builder so the drain can /// `give_fueled` each emitted update and yield to timely under large snapshot @@ -442,7 +466,55 @@ where source_config.source_statistics.clone(), ); match flavor { - UpsertStashFlavor::Chunked => { + UpsertStashFlavor::Payload => { + let pool = mz_timely_util::columnar::chunk::spill_pool() + .or_else(mz_timely_util::pool_config::global_pool) + .expect("payload upsert requires a buffer pool"); + let store = payload::store(pool); + let (encoded, token) = payload::encode_feedback(encoded, store.clone()); + let persist_arranged = arrange_core::< + _, + _, + payload::FeedbackChunker, + ChunkBatcher>, + ChunkBuilder>, + payload::FeedbackSpine, + >(encoded, Pipeline, "Persist payload feedback"); + let mut persist_token = persist_token.unwrap_or_default(); + persist_token.push(token); + build_upsert_operator::( + input, + resume_upper, + persist_arranged, + Some(persist_token), + upsert_metrics, + source_config, + true, + Some(store), + ) + } + + UpsertStashFlavor::Resumable { async_reads } => { + let (persist_arranged, token) = mz_timely_util::columnar::chunk::asynchronous::arrange( + encoded, + merge_read_budget(), + "Persist resumable feedback", + ); + let mut persist_token = persist_token.unwrap_or_default(); + persist_token.push(token); + build_upsert_operator::( + input, + resume_upper, + persist_arranged, + Some(persist_token), + upsert_metrics, + source_config, + async_reads, + None, + ) + } + + UpsertStashFlavor::Chunked { async_reads } => { // Chains and sealed batches alike are `FeedbackChunk`s whose // bodies spill to the buffer pool, behind the same process spill // gate as the source stash. @@ -461,6 +533,8 @@ where persist_token, upsert_metrics, source_config, + async_reads, + None, ) } UpsertStashFlavor::Paged => { @@ -484,6 +558,8 @@ where persist_token, upsert_metrics, source_config, + false, + None, ) } } @@ -568,6 +644,8 @@ fn build_upsert_operator<'scope, A, T, FromTime>( persist_token: Option>, upsert_metrics: UpsertMetrics, source_config: crate::source::SourceExportCreationConfig, + async_reads: bool, + payload_store: Option, ) -> ( VecCollection<'scope, T, Result, Diff>, StreamVec<'scope, T, (Option, HealthStatusUpdate)>, @@ -576,6 +654,7 @@ fn build_upsert_operator<'scope, A, T, FromTime>( ) where A: UpsertStashArm, + for<'a> columnar::Ref<'a, A::Value>: Copy + Ord, T: Timestamp + TotalOrder + Sync, T: Refines + differential_dataflow::lattice::Lattice, T: columnation::Columnation, @@ -629,13 +708,13 @@ where // pushed in, bounding memory to O(unique key-time pairs) even during // large initial snapshots. How cold stash state leaves RSS is // flavor-specific; see the arm impls. - let mut batcher = A::new_batcher(); + let mut batcher = A::new_batcher(payload_store); // The chunker sorts and consolidates raw input into the `Column` chunks // the batcher consumes. - let mut chunker: UpsertChunker = Default::default(); + let mut chunker: UpsertChunker = Default::default(); // Scratch buffer for accumulating source events before flushing to // the batcher. Drained on each iteration via the chunker. - let mut push_buffer: Vec> = Vec::new(); + let mut push_buffer: Vec> = Vec::new(); // Capability held at the minimum time of any buffered data. When // Some, the operator may still produce output; when None, the @@ -660,7 +739,13 @@ where tokio::select! { _ = input.ready() => {} _ = persist_wakeup.ready() => { - while persist_wakeup.next_sync().is_some() {} + while let Some(event) = persist_wakeup.next_sync() { + if let AsyncEvent::Data(_, batches) = event { + for batch in batches { + mz_timely_util::columnar::chunk::metrics::record_batch(batch.len()); + } + } + } } } @@ -679,7 +764,8 @@ where { continue; } - let value = value.as_ref().map(upsert_value_to_row); + let value = + A::encode(value.as_ref().map(upsert_value_to_row), &mut batcher); let from_time = from_time.upsert_order(); push_buffer.push((key, ts, UpsertDiff { from_time, value })); pushed_any = true; @@ -706,7 +792,7 @@ where // Flush buffered events through the chunker into the batcher. This // triggers the chunker + geometric chain merging, which consolidates // entries for the same (key, time) via the UpsertDiff Semigroup. - A::flush(&mut push_buffer, &mut chunker, &mut batcher); + A::flush(&mut push_buffer, &mut chunker, &mut batcher).await; // Step 2: Read persist frontier. // The persist probe tells us which output times have been @@ -789,9 +875,9 @@ where // Step 1 already consolidated `push_buffer` through the chunker // (which readies a complete chunk per `push_into`), so the // chunker holds nothing pending here and we can seal directly. - let (sealed, _description) = batcher.seal(input_upper.clone()); + let sealed = A::seal(&mut batcher, input_upper.clone()).await; // Frontier of data remaining in the batcher (ts >= input_upper). - let remaining_frontier = batcher.frontier().to_owned(); + let remaining_frontier = A::frontier(&mut batcher); let mut ineligible = Vec::new(); // The drain emits eligible output directly through @@ -806,6 +892,8 @@ where &mut persist_trace, source_config.worker_id, source_config.id, + async_reads, + &mut batcher, ) .await; @@ -830,7 +918,7 @@ where // remaining data: either entries still in the batcher (above // input_upper) or ineligible entries being pushed back. let min_ineligible_ts = ineligible.iter().map(|(_, ts, _)| ts).min().cloned(); - A::flush(&mut ineligible, &mut chunker, &mut batcher); + A::flush(&mut ineligible, &mut chunker, &mut batcher).await; // `Option::min` alone would be wrong here, `None` sorts low. // Chain the candidates and take the min over present ones. @@ -878,54 +966,76 @@ where O: columnar::Columnar + Default + Ord + Clone + Send + Sync + 'static, for<'a> columnar::Ref<'a, O>: Ord + Copy, { + /// Representation of a value in source metadata. + type Value: columnar::Columnar + Default + Clone; + + /// Prepare a source value for offset consolidation. + fn encode(value: Option, batcher: &mut Self::Batcher) -> Option; + /// Release temporary owners after all buffered updates have entered chunks. + fn end_flush(_batcher: &mut Self::Batcher) {} + /// The feedback arrangement's spine. `'static` because the operator /// future owns a trace agent for it. type Spine: TraceReader