{"status":"ok","message-type":"work","message-version":"1.0.0","message":{"indexed":{"date-parts":[[2026,7,3]],"date-time":"2026-07-03T16:28:27Z","timestamp":1783096107684,"version":"3.54.6"},"reference-count":35,"publisher":"Association for Computing Machinery (ACM)","issue":"4","content-domain":{"domain":[],"crossmark-restriction":false},"short-container-title":["Proc. VLDB Endow."],"published-print":{"date-parts":[[2021,12]]},"abstract":"<jats:p>\n            The cost of big-data query execution is dominated by stateful operators. These include\n            <jats:italic>sort<\/jats:italic>\n            and\n            <jats:italic>hash-aggregate<\/jats:italic>\n            that typically materialize intermediate data in memory, and\n            <jats:italic>exchange<\/jats:italic>\n            that materializes data to disk and transfers data over the network. In this paper we focus on several query optimization techniques that reduce the cost of these operators. First, we introduce a novel exchange placement algorithm that improves the state-of-the-art and significantly reduces the amount of data exchanged. The algorithm simultaneously minimizes the number of exchanges required and maximizes computation reuse via multi-consumer exchanges. Second, we introduce three partial push-down optimizations that push down partial computation derived from existing operators (\n            <jats:italic>group-bys<\/jats:italic>\n            ,\n            <jats:italic>intersections<\/jats:italic>\n            and\n            <jats:italic>joins<\/jats:italic>\n            ) below these stateful operators. While these optimizations are generically applicable we find that two of these optimizations (\n            <jats:italic>partial aggregate and partial semi-join push-down<\/jats:italic>\n            ) are only beneficial in the scale-out setting where\n            <jats:italic>exchanges<\/jats:italic>\n            are a bottleneck. We propose novel extensions to existing literature to perform more aggressive partial push-downs than the state-of-the-art and also specialize them to the big-data setting. Finally we propose peephole optimizations that specialize the implementation of stateful operators to their input parameters. All our optimizations are implemented in the spark engine that powers azure synapse. We evaluate their impact on TPCDS and demonstrate that they make our engine 1.8X faster than Apache Spark 3.0.1.\n          <\/jats:p>","DOI":"10.14778\/3503585.3503601","type":"journal-article","created":{"date-parts":[[2022,4,14]],"date-time":"2022-04-14T22:18:07Z","timestamp":1649974687000},"page":"936-948","source":"Crossref","is-referenced-by-count":12,"title":["New query optimization techniques in the Spark engine of Azure synapse"],"prefix":"10.14778","volume":"15","author":[{"given":"Abhishek","family":"Modi","sequence":"first","affiliation":[{"name":"Microsoft, India"}],"role":[{"vocabulary":"crossref","role":"author"}]},{"given":"Kaushik","family":"Rajan","sequence":"additional","affiliation":[{"name":"Microsoft Research, India"}],"role":[{"vocabulary":"crossref","role":"author"}]},{"given":"Srinivas","family":"Thimmaiah","sequence":"additional","affiliation":[{"name":"Microsoft, India"}],"role":[{"vocabulary":"crossref","role":"author"}]},{"given":"Prakhar","family":"Jain","sequence":"additional","affiliation":[{"name":"Databricks"}],"role":[{"vocabulary":"crossref","role":"author"}]},{"given":"Swinky","family":"Mann","sequence":"additional","affiliation":[{"name":"Microsoft, India"}],"role":[{"vocabulary":"crossref","role":"author"}]},{"given":"Ayushi","family":"Agarwal","sequence":"additional","affiliation":[{"name":"Microsoft, India"}],"role":[{"vocabulary":"crossref","role":"author"}]},{"given":"Ajith","family":"Shetty","sequence":"additional","affiliation":[{"name":"Microsoft, India"}],"role":[{"vocabulary":"crossref","role":"author"}]},{"given":"Shahid K","family":"I","sequence":"additional","affiliation":[{"name":"Microsoft, India"}],"role":[{"vocabulary":"crossref","role":"author"}]},{"given":"Ashit","family":"Gosalia","sequence":"additional","affiliation":[{"name":"Microsoft"}],"role":[{"vocabulary":"crossref","role":"author"}]},{"given":"Partho","family":"Sarthi","sequence":"additional","affiliation":[{"name":"University of Wisconsin-Madison"}],"role":[{"vocabulary":"crossref","role":"author"}]}],"member":"320","published-online":{"date-parts":[[2022,4,14]]},"reference":[{"key":"e_1_2_1_1_1","unstructured":"Spark SQL Aggregate Rewrite Rule. https:\/\/github.com\/apache\/spark\/blob\/master\/sql\/core\/src\/main\/scala\/org\/apache\/spark\/sql\/execution\/aggregate\/AggUtils.scala.  Spark SQL Aggregate Rewrite Rule. https:\/\/github.com\/apache\/spark\/blob\/master\/sql\/core\/src\/main\/scala\/org\/apache\/spark\/sql\/execution\/aggregate\/AggUtils.scala."},{"key":"e_1_2_1_2_1","unstructured":"Spark SQL Expand Operator. https:\/\/github.com\/apache\/spark\/blob\/master\/sql\/catalyst\/src\/main\/scala\/org\/apache\/spark\/sql\/catalyst\/plans\/logical\/basicLogicalOperators.scala.  Spark SQL Expand Operator. https:\/\/github.com\/apache\/spark\/blob\/master\/sql\/catalyst\/src\/main\/scala\/org\/apache\/spark\/sql\/catalyst\/plans\/logical\/basicLogicalOperators.scala."},{"key":"e_1_2_1_3_1","unstructured":"Spark SQL HashAggregate Operator. https:\/\/github.com\/apache\/spark\/blob\/master\/sql\/core\/src\/main\/scala\/org\/apache\/spark\/sql\/execution\/aggregate\/HashAggregateExec.scala.  Spark SQL HashAggregate Operator. https:\/\/github.com\/apache\/spark\/blob\/master\/sql\/core\/src\/main\/scala\/org\/apache\/spark\/sql\/execution\/aggregate\/HashAggregateExec.scala."},{"key":"e_1_2_1_4_1","unstructured":"Spark SQL Set Operators. https:\/\/spark.apache.org\/docs\/latest\/sql-ref-syntax-qry-select-setops.html.  Spark SQL Set Operators. https:\/\/spark.apache.org\/docs\/latest\/sql-ref-syntax-qry-select-setops.html."},{"key":"e_1_2_1_5_1","volume-title":"https:\/\/databricks.com\/blog\/2014\/10\/10\/spark-petabyte-sort.html","author":"Fastest Open Source Apache Spark","year":"2014","unstructured":"Apache Spark the Fastest Open Source Engine for Sorting a Petabyte . https:\/\/databricks.com\/blog\/2014\/10\/10\/spark-petabyte-sort.html , 2014 . Apache Spark the Fastest Open Source Engine for Sorting a Petabyte. https:\/\/databricks.com\/blog\/2014\/10\/10\/spark-petabyte-sort.html, 2014."},{"key":"e_1_2_1_6_1","doi-asserted-by":"publisher","DOI":"10.1145\/2723372.2742797"},{"key":"e_1_2_1_7_1","doi-asserted-by":"publisher","DOI":"10.1145\/320455.320457"},{"key":"e_1_2_1_8_1","doi-asserted-by":"publisher","DOI":"10.1145\/362686.362692"},{"key":"e_1_2_1_9_1","doi-asserted-by":"publisher","DOI":"10.1145\/275487.275492"},{"key":"e_1_2_1_10_1","doi-asserted-by":"publisher","DOI":"10.5555\/645920.672834"},{"key":"e_1_2_1_11_1","doi-asserted-by":"publisher","DOI":"10.1109\/71.159044"},{"key":"e_1_2_1_12_1","first-page":"02","article-title":"On applying hash filters to improving the execution of multi-join queries","volume":"6","author":"Yu Philip","year":"2000","unstructured":"Ming-syan Chen, Hui-i Hsiao, and Philip Yu . On applying hash filters to improving the execution of multi-join queries . The VLDB Journal The International Journal on Very Large Data Bases , 6 , 02 2000 . Ming-syan Chen, Hui-i Hsiao, and Philip Yu. On applying hash filters to improving the execution of multi-join queries. The VLDB Journal The International Journal on Very Large Data Bases, 6, 02 2000.","journal-title":"The VLDB Journal The International Journal on Very Large Data Bases"},{"key":"e_1_2_1_13_1","doi-asserted-by":"publisher","DOI":"10.5555\/645919.672666"},{"key":"e_1_2_1_14_1","doi-asserted-by":"publisher","DOI":"10.1145\/2882903.2903741"},{"key":"e_1_2_1_15_1","doi-asserted-by":"publisher","DOI":"10.1145\/3318464.3389769"},{"key":"e_1_2_1_16_1","doi-asserted-by":"publisher","DOI":"10.1145\/2674005.2674994"},{"key":"e_1_2_1_17_1","doi-asserted-by":"publisher","DOI":"10.14778\/2168651.2168654"},{"key":"e_1_2_1_18_1","doi-asserted-by":"publisher","DOI":"10.1145\/152610.152611"},{"key":"e_1_2_1_19_1","first-page":"18","article-title":"The cascades framework for query optimization","author":"Graefe Goetz","year":"1995","unstructured":"Goetz Graefe . The cascades framework for query optimization . Data Engineering Bulletin , 18 , 1995 . Goetz Graefe. The cascades framework for query optimization. Data Engineering Bulletin, 18, 1995.","journal-title":"Data Engineering Bulletin"},{"key":"e_1_2_1_20_1","doi-asserted-by":"publisher","DOI":"10.5555\/645921.673150"},{"key":"e_1_2_1_21_1","doi-asserted-by":"publisher","DOI":"10.14778\/3192965.3192971"},{"key":"e_1_2_1_22_1","doi-asserted-by":"publisher","DOI":"10.1109\/ICDE.2002.994787"},{"key":"e_1_2_1_23_1","volume-title":"Efficiently compiling efficient query plans for modern hardware. PVLDB, 4(9)","author":"Neumann Thomas","year":"2011","unstructured":"Thomas Neumann . Efficiently compiling efficient query plans for modern hardware. PVLDB, 4(9) , 2011 . Thomas Neumann. Efficiently compiling efficient query plans for modern hardware. PVLDB, 4(9), 2011."},{"key":"e_1_2_1_24_1","first-page":"293","volume-title":"12th USENIX Symposium on Networked Systems Design and Implementation (NSDI 15)","author":"Ousterhout Kay","year":"2015","unstructured":"Kay Ousterhout , Ryan Rasti , Sylvia Ratnasamy , Scott Shenker , and Byung-Gon Chun . Making sense of performance in data analytics frameworks . In 12th USENIX Symposium on Networked Systems Design and Implementation (NSDI 15) , pages 293 -- 307 , Oakland, CA , May 2015 . USENIX Association. Kay Ousterhout, Ryan Rasti, Sylvia Ratnasamy, Scott Shenker, and Byung-Gon Chun. Making sense of performance in data analytics frameworks. In 12th USENIX Symposium on Networked Systems Design and Implementation (NSDI 15), pages 293--307, Oakland, CA, May 2015. USENIX Association."},{"key":"e_1_2_1_25_1","doi-asserted-by":"publisher","DOI":"10.5555\/1070432.1070548"},{"key":"e_1_2_1_26_1","doi-asserted-by":"publisher","DOI":"10.14778\/3339490.3339495"},{"key":"e_1_2_1_27_1","doi-asserted-by":"publisher","DOI":"10.1109\/ICDE.2019.00196"},{"key":"e_1_2_1_28_1","doi-asserted-by":"publisher","DOI":"10.14778\/3415478.3415558"},{"key":"e_1_2_1_29_1","doi-asserted-by":"publisher","DOI":"10.14778\/1687553.1687609"},{"key":"e_1_2_1_30_1","first-page":"345","volume-title":"Proceedings of the 21th International Conference on Very Large Data Bases, VLDB '95","author":"Weipeng","year":"1995","unstructured":"Weipeng P. Yan and Per-\u00c5ke Larson. Eager aggregation and lazy aggregation . In Proceedings of the 21th International Conference on Very Large Data Bases, VLDB '95 , page 345 -- 357 , San Francisco, CA, USA , 1995 . Morgan Kaufmann Publishers Inc. Weipeng P. Yan and Per-\u00c5ke Larson. Eager aggregation and lazy aggregation. In Proceedings of the 21th International Conference on Very Large Data Bases, VLDB '95, page 345--357, San Francisco, CA, USA, 1995. Morgan Kaufmann Publishers Inc."},{"key":"e_1_2_1_31_1","doi-asserted-by":"publisher","DOI":"10.1145\/1629575.1629600"},{"key":"e_1_2_1_32_1","doi-asserted-by":"publisher","DOI":"10.1145\/3190508.3190534"},{"key":"e_1_2_1_33_1","first-page":"295","volume-title":"9th USENIX Symposium on Networked Systems Design and Implementation (NSDI 12)","author":"Zhang Jiaxing","year":"2012","unstructured":"Jiaxing Zhang , Hucheng Zhou , Rishan Chen , Xuepeng Fan , Zhenyu Guo , Haoxiang Lin , Jack Y. Li , Wei Lin , Jingren Zhou , and Lidong Zhou . Optimizing data shuffling in data-parallel computation by understanding user-defined functions . In 9th USENIX Symposium on Networked Systems Design and Implementation (NSDI 12) , pages 295 -- 308 , San Jose, CA , April 2012 . USENIX Association. Jiaxing Zhang, Hucheng Zhou, Rishan Chen, Xuepeng Fan, Zhenyu Guo, Haoxiang Lin, Jack Y. Li, Wei Lin, Jingren Zhou, and Lidong Zhou. Optimizing data shuffling in data-parallel computation by understanding user-defined functions. In 9th USENIX Symposium on Networked Systems Design and Implementation (NSDI 12), pages 295--308, San Jose, CA, April 2012. USENIX Association."},{"key":"e_1_2_1_34_1","doi-asserted-by":"publisher","DOI":"10.1109\/ICDE.2010.5447802"},{"key":"e_1_2_1_35_1","doi-asserted-by":"publisher","DOI":"10.1007\/s00778-012-0280-z"}],"container-title":["Proceedings of the VLDB Endowment"],"original-title":[],"language":"en","link":[{"URL":"https:\/\/dl.acm.org\/doi\/pdf\/10.14778\/3503585.3503601","content-type":"unspecified","content-version":"vor","intended-application":"similarity-checking"}],"deposited":{"date-parts":[[2022,12,28]],"date-time":"2022-12-28T10:31:13Z","timestamp":1672223473000},"score":1,"resource":{"primary":{"URL":"https:\/\/dl.acm.org\/doi\/10.14778\/3503585.3503601"}},"subtitle":[],"short-title":[],"issued":{"date-parts":[[2021,12]]},"references-count":35,"journal-issue":{"issue":"4","published-print":{"date-parts":[[2021,12]]}},"alternative-id":["10.14778\/3503585.3503601"],"URL":"https:\/\/doi.org\/10.14778\/3503585.3503601","relation":{},"ISSN":["2150-8097"],"issn-type":[{"value":"2150-8097","type":"print"}],"subject":[],"published":{"date-parts":[[2021,12]]}}}