SQL 파이프 구문
SQL 파이프 구문 (SQL Pipe Syntax)
Apache Spark의 SQL pipe syntax는 파이프 문자 |>를 사용해 연산자들을 조합해 쿼리를 구성하는 방법이에요. 이 구문이 어떻게 동작하는지, 프로젝션·집계·조인 등의 각 연산자 의미와 예제를 정리해서 알아볼게요.
출처: 문서
본문
구문 (Syntax)
개요 (Overview)
Apache Spark는 연산자(operator)들의 조합으로 쿼리를 구성할 수 있게 해주는 SQL pipe syntax를 지원해요.
- 모든 쿼리는 파이프 문자
|>로 구분되는 0개 이상의 파이프 연산자를 접미사로 가질 수 있어요. - 각 파이프 연산자는 아래 표에 설명된 것처럼 하나 이상의 SQL 키워드로 시작하고 그 뒤에 고유한 문법을 가져요.
- 대부분의 연산자는 표준 SQL 절을 위한 기존 문법을 재사용해요.
- 연산자는 어떤 순서로든, 몇 번이든 적용할 수 있어요.
FROM <tableName>은 이제 TABLE <tableName>과 동일하게 동작하는, 지원되는 독립형(standalone) 쿼리가 됐어요. 이는 연결된(체인) 파이프 SQL 쿼리를 시작하기에 편리한 출발점을 제공해요. 물론 유효한 Spark SQL 쿼리라면 무엇이든 그 끝에 하나 이상의 파이프 연산자를 추가할 수 있으며, 여기에 작성된 것과 동일한 일관된 동작을 보여요.
지원되는 모든 연산자와 그 의미론의 전체 목록은 이 문서 끝의 표를 참고하세요.
예제 (Example)
예를 들어, 이것은 TPC-H 벤치마크의 쿼리 13이에요:
SELECT c_count, COUNT(*) AS custdist
FROM
(SELECT c_custkey, COUNT(o_orderkey) c_count
FROM customer
LEFT OUTER JOIN orders ON c_custkey = o_custkey
AND o_comment NOT LIKE '%unusual%packages%' GROUP BY c_custkey
) AS c_orders
GROUP BY c_count
ORDER BY custdist DESC, c_count DESC;
같은 로직을 SQL 파이프 연산자로 작성하면 이렇게 표현해요:
FROM customer
|> LEFT OUTER JOIN orders ON c_custkey = o_custkey
AND o_comment NOT LIKE '%unusual%packages%'
|> AGGREGATE COUNT(o_orderkey) c_count
GROUP BY c_custkey
|> AGGREGATE COUNT(*) AS custdist
GROUP BY c_count
|> ORDER BY custdist DESC, c_count DESC;
소스 테이블 (Source Tables)
SQL pipe syntax로 새 쿼리를 시작하려면 FROM <tableName> 또는 TABLE <tableName> 절을 사용해요. 이는 소스 테이블의 모든 행을 포함하는 관계를 생성해요. 그런 다음 이 절의 끝에 하나 이상의 파이프 연산자를 추가해 추가 변환을 수행해요.
프로젝션 (Projections)
SQL pipe syntax는 표현식을 평가하는 구성 가능한(composable) 방법을 지원해요. 이러한 프로젝션 기능의 큰 장점은 이전 표현식에 기반한 새 표현식을 증분 방식으로 계산할 수 있다는 점이에요. 각 연산자는 연산자가 나타나는 순서와 무관하게 입력 테이블에 독립적으로 적용되므로, 여기서 측면 컬럼 참조(lateral column references)가 필요 없어요. 계산된 각 컬럼은 다음 연산자와 함께 사용할 수 있게 됩니다.
SELECT는 제공된 표현식을 평가해 새 테이블을 생성해요. 필요에 따라 DISTINCT와 *를 사용할 수 있어요. 이는 일반 Spark SQL에서 테이블 하위 쿼리의 가장 바깥쪽 SELECT처럼 동작해요.
EXTEND는 제공된 표현식을 평가해 입력 테이블에 새 컬럼을 추가해요. 이는 테이블 별칭도 보존해요. 일반 Spark SQL의 SELECT *, new_column처럼 동작해요.
DROP는 입력 테이블에서 컬럼을 제거해요. 일반 Spark SQL의 SELECT * EXCEPT (column)과 비슷해요.
SET는 입력 테이블의 컬럼 값을 교체해요. 일반 Spark SQL의 SELECT * REPLACE (expression AS column)과 비슷해요.
AS는 입력 테이블을 전달하고 각 행에 새 별칭을 도입해요.
집계 (Aggregations)
일반적으로 SQL pipe syntax에서의 집계는 일반 Spark SQL과 달리 다르게 이루어져요.
전체 테이블 집계를 수행하려면 평가할 집계 표현식 목록과 함께 AGGREGATE 연산자를 사용해요. 이는 출력 테이블에 단일 행 하나를 반환해요.
그룹화와 함께 집계를 수행하려면 GROUP BY 절이 있는 AGGREGATE 연산자를 사용해요. 이는 그룹화 표현식 값의 각 고유 조합에 대해 하나의 행을 반환해요. 출력 테이블은 평가된 그룹화 표현식 다음에 평가된 집계 함수를 포함해요. 그룹화 표현식은 향후 연산자에서 참조하기 위한 목적으로 별칭(alias) 할당을 지원해요. 이렇게 하면 AGGREGATE가 둘 다 수행하는 단일 연산자이므로, GROUP BY와 SELECT 사이에 전체 표현식을 반복할 필요가 없어요.
기타 변환 (Other Transformations)
나머지 연산자는 필터링, 조인, 정렬, 샘플링, 집합 연산과 같은 다른 변환에 사용돼요. 이 연산자들은 일반적으로 아래 표에 설명된 대로 일반 Spark SQL에서와 같은 방식으로 동작해요.
독립성과 상호 운용성 (Independence and Interoperability)
SQL pipe syntax는 기존 SQL 쿼리에 대한 하위 호환성 문제 없이 Spark에서 동작해요. 어떤 쿼리든 일반 Spark SQL, pipe syntax, 또는 그 둘의 조합으로 작성할 수 있어요. 그 결과 다음 불변식(invariants)이 항상 성립해요:
- 각 파이프 연산자는 입력 테이블을 받아, 그것이 어떻게 계산됐는지와 무관하게 그 행들에 대해 같은 방식으로 동작해요.
- N개의 SQL 파이프 연산자로 구성된 유효한 체인에 대해, 처음 M <= N개의 연산자의 어떤 부분집합도 항상 유효한 쿼리를 나타내요.
이 속성은 Jupyter 노트북 같은 SQL 편집기의 "run highlighted text" 기능처럼 행들의 부분집합을 선택하는 등의 내관(introspection)·디버깅에 유용할 수 있어요.
- 일반 Spark SQL로 작성된 유효한 쿼리라면 무엇이든 파이프 연산자를 추가할 수 있어요.
파이프 syntax 쿼리를 시작하는 표준적인 방법은 FROM <tableName> 절이에요.
참고로 이는 유효한 독립형 쿼리이며, 일반성을 잃지 않고 다른 Spark SQL 쿼리로 대체될 수 있어요.
- 테이블 하위 쿼리는 일반 Spark SQL syntax 또는 pipe syntax 중 하나로 작성할 수 있어요.
이들은 어느 syntax로 작성된 내포하는 쿼리 안에도 나타날 수 있어요.
- 뷰 및 DDL·DML 명령과 같은 다른 Spark SQL 문은 어느 syntax로 작성된 쿼리든 포함할 수 있어요.
지원되는 연산자 (Supported Operators)
| 연산자 (Operator) | 출력 행 (Output rows) |
|---|---|
FROM or TABLE |
소스 테이블의 모든 출력 행을 수정하지 않고 반환해요. |
SELECT |
입력 테이블의 각 행에 대해 제공된 표현식을 평가해요. |
EXTEND |
각 입력 행에 대해 지정된 표현식을 평가해 입력 테이블에 새 컬럼을 추가해요. |
SET |
제공된 표현식을 평가한 결과로 컬럼을 교체해 입력 테이블의 컬럼을 갱신해요. |
DROP |
이름으로 입력 테이블의 컬럼을 제거해요. |
AS |
입력 테이블과 같은 행·컬럼 이름을 유지하지만 새 테이블 별칭을 가져요. |
WHERE |
조건을 통과하는 입력 행의 부분집합을 반환해요. |
LIMIT |
(있을 경우) 순서를 유지하며 지정된 수의 입력 행을 반환해요. |
AGGREGATE |
그룹화 유무에 관계없이 집계를 수행해요. |
JOIN |
두 입력의 행을 결합해, 입력 테이블과 테이블 인자의 필터링된 교차곱을 반환해요. |
ORDER BY |
지시된 대로 정렬한 후 입력 행을 반환해요. |
UNION ALL |
입력 테이블과 다른 테이블 인자(들)의 결합된 행에 대해 union 또는 기타 집합 연산을 수행해요. |
TABLESAMPLE |
제공된 샘플링 알고리즘이 선택한 행의 부분집합을 반환해요. |
PIVOT |
입력 행이 컬럼이 되도록 피벗된 새 테이블을 반환해요. |
UNPIVOT |
입력 컬럼이 행이 되도록 피벗된 새 테이블을 반환해요. |
이 표는 지원되는 각 파이프 연산자와 그들이 생성하는 출력 행을 설명해요. 각 연산자는 |> 기호 앞의 쿼리가 생성한 행으로 구성된 입력 관계를 받아들인다는 점을 기억하세요.
FROM 또는 TABLE (FROM or TABLE)
FROM <tableName>
TABLE <tableName>
소스 테이블의 모든 출력 행을 수정하지 않고 반환해요.
예를 들어:
CREATE TABLE t AS VALUES (1, 2), (3, 4) AS t(a, b);
TABLE t;
+---+---+
| a| b|
+---+---+
| 1| 2|
| 3| 4|
+---+---+
SELECT
|> SELECT <expr> [[AS] alias], ...
입력 테이블의 각 행에 대해 제공된 표현식을 평가해요.
일반적으로 SQL pipe syntax에서 이 연산자가 항상 필요한 것은 아니에요. 쿼리의 끝이나 끝 근처에서 표현식을 평가하거나 출력 컬럼 목록을 지정하는 데 사용할 수 있어요.
최종 쿼리 결과는 항상 마지막 파이프 연산자가 반환한 컬럼들로 구성되므로, 이 SELECT 연산자가 나타나지 않으면 출력에는 전체 행의 모든 컬럼이 포함돼요. 이 동작은 표준 SQL syntax의 SELECT *과 비슷해요.
필요에 따라 DISTINCT와 *을 사용할 수 있어요.
이는 일반 Spark SQL에서 테이블 하위 쿼리의 가장 바깥쪽 SELECT처럼 동작해요.
SELECT 목록에서 윈도우 함수(window functions)도 지원돼요. 이를 사용하려면 OVER 절이 제공되어야 해요. 윈도우 명세를 WINDOW 절에 제공할 수 있어요.
이 연산자에서는 집계 함수가 지원되지 않아요. 집계를 수행하려면 대신 AGGREGATE 연산자를 사용하세요.
예를 들어:
CREATE TABLE t AS VALUES (0), (1) AS t(col);
FROM t
|> SELECT col * 2 AS result;
+------+
|result|
+------+
| 0|
| 2|
+------+
EXTEND
|> EXTEND <expr> [[AS] alias], ...
각 입력 행에 대해 지정된 표현식을 평가해 입력 테이블에 새 컬럼을 추가해요.
EXTEND 연산 후에는 최상위 컬럼 이름이 갱신되지만 테이블 별칭은 여전히 원래 행 값을 가리켜요 (예: 두 테이블 lhs와 rhs 사이의 inner join 후에 EXTEND를 수행하고 나서 SELECT lhs.col, rhs.col을 하는 경우).
예를 들어:
VALUES (0), (1) tab(col)
|> EXTEND col * 2 AS result;
+---+------+
|col|result|
+---+------+
| 0| 0|
| 1| 2|
+---+------+
SET
|> SET <column> = <expression>, ...
제공된 표현식을 평가한 결과로 컬럼을 교체해 입력 테이블의 컬럼을 갱신해요. 각 컬럼 참조는 입력 테이블에 정확히 한 번만 나타나야 해요.
일반 Spark SQL의 SELECT * EXCEPT (column), <expression> AS column과 비슷해요.
단일 SET 절에서 여러 할당을 수행할 수 있어요. 각 할당은 이전 할당의 결과를 참조할 수 있어요.
할당 후에는 최상위 컬럼 이름이 갱신되지만 테이블 별칭은 여전히 원래 행 값을 가리켜요 (예: 두 테이블 lhs와 rhs 사이의 inner join 후에 SET을 수행하고 나서 SELECT lhs.col, rhs.col을 하는 경우).
예를 들어:
VALUES (0), (1) tab(col)
|> SET col = col * 2;
+---+
|col|
+---+
| 0|
| 2|
+---+
VALUES (0), (1) tab(col)
|> SET col = col * 2;
+---+
|col|
+---+
| 0|
| 2|
+---+
DROP
|> DROP <column>, ...
이름으로 입력 테이블의 컬럼을 제거해요. 각 컬럼 참조는 입력 테이블에 정확히 한 번만 나타나야 해요.
일반 Spark SQL의 SELECT * EXCEPT (column)과 비슷해요.
DROP 연산 후에는 최상위 컬럼 이름이 갱신되지만 테이블 별칭은 여전히 원래 행 값을 가리켜요 (예: 두 테이블 lhs와 rhs 사이의 inner join 후에 DROP을 수행하고 나서 SELECT lhs.col, rhs.col을 하는 경우).
예를 들어:
VALUES (0, 1) tab(col1, col2)
|> DROP col1;
+----+
|col2|
+----+
| 1|
+----+
AS
|> AS <alias>
입력 테이블과 같은 행·컬럼 이름을 유지하지만 새 테이블 별칭을 가져요.
이 연산자는 입력 테이블에 새 별칭을 도입하는 데 유용하며, 이후 연산자에서 참조할 수 있어요. 테이블의 기존 별칭은 새 별칭으로 교체돼요.
SELECT나 EXTEND로 새 컬럼을 추가한 후, 또는 AGGREGATE로 집계를 수행한 후에 이 연산자를 사용하는 것이 유용해요. 이는 후속 JOIN 연산자에서 컬럼을 참조하는 과정을 단순화하고 더 읽기 쉬운 쿼리를 허용해요.
예를 들어:
VALUES (0, 1) tab(col1, col2)
|> AS new_tab
|> SELECT col1 + col2 FROM new_tab;
+-----------+
|col1 + col2|
+-----------+
| 1|
+-----------+
WHERE
|> WHERE <condition>
조건을 통과하는 입력 행의 부분집합을 반환해요.
이 연산자는 어디든 나타날 수 있으므로 별도의 HAVING이나 QUALIFY 구문이 필요 없어요.
예를 들어:
VALUES (0), (1) tab(col)
|> WHERE col = 1;
+---+
|col|
+---+
| 1|
+---+
LIMIT
|> [LIMIT <n>] [OFFSET <m>]
(있을 경우) 순서를 유지하며 지정된 수의 입력 행을 반환해요.
LIMIT와 OFFSET은 함께 지원돼요. LIMIT 절은 OFFSET 절 없이도 사용할 수 있고, OFFSET 절도 LIMIT 절 없이 사용할 수 있어요.
예를 들어:
VALUES (0), (0) tab(col)
|> LIMIT 1;
+---+
|col|
+---+
| 0|
+---+
AGGREGATE
-- Full-table aggregation
|> AGGREGATE <agg_expr> [[AS] alias], ...
-- Aggregation with grouping
|> AGGREGATE [<agg_expr> [[AS] alias], ...] GROUP BY <grouping_expr> [AS alias], ...
그룹화된 행 또는 전체 입력 테이블에 걸쳐 집계를 수행해요.
GROUP BY 절이 없으면 전체 테이블 집계를 수행해 각 집계 표현식에 대해 하나의 컬럼을 가진 결과 행 하나를 반환해요. 그렇지 않으면 그룹화와 함께 집계를 수행해, 그룹당 하나의 행을 반환해요. 그룹화 표현식에 직접 별칭을 할당할 수 있어요.
이 연산자의 출력 컬럼 목록은 (있을 경우) 그룹화 컬럼이 먼저, 그 다음 집계 컬럼이 뒤따라요.
각 <agg_expr> 표현식은 COUNT, SUM, AVG, MIN 또는 Spark SQL이 지원하는 다른 집계 함수와 같은 표준 집계 함수를 포함할 수 있어요. 집계 함수 위나 아래에 추가 표현식이 나타날 수 있어요 (예: MIN(FLOOR(col)) + 1). 각 <agg_expr> 표현식은 최소한 하나의 집계 함수를 포함해야 해요 (그렇지 않으면 쿼리가 오류를 반환해요). 각 <agg_expr> 표현식은 AS <alias>로 컬럼 별칭을 포함할 수 있고, 집계 함수를 적용하기 전에 중복 값을 제거하기 위한 DISTINCT 키워드도 포함할 수 있어요 (예: COUNT(DISTINCT col)).
GROUP BY 절이 있으면 임의 개수의 그룹화 표현식을 포함할 수 있고, 각 <agg_expr> 표현식은 그룹화 표현식 값의 각 고유 조합에 대해 평가돼요. 출력 테이블은 평가된 그룹화 표현식 다음에 평가된 집계 함수를 포함해요. GROUP BY 표현식은 1-기반 순서수(ordinals)를 포함할 수 있어요. 그러한 순서수가 동반된 SELECT 절의 표현식을 참조하는 일반 SQL과 달리, SQL pipe syntax에서는 이전 연산자가 생성한 관계의 컬럼을 참조해요. 예를 들어 TABLE t |> AGGREGATE COUNT(*) GROUP BY 2에서 우리는 입력 테이블 t의 두 번째 컬럼을 참조해요.
AGGREGATE 연산자가 출력에 평가된 그룹화 표현식을 자동으로 포함하므로 GROUP BY와 SELECT 사이에 전체 표현식을 반복할 필요가 없어요. 마찬가지로 AGGREGATE 연산자 후에는 보통 후속 SELECT 연산자를 발행할 필요가 없는데, AGGREGATE가 그룹화 컬럼과 집계 컬럼을 단일 단계로 반환하기 때문이에요.
예를 들어:
-- Full-table aggregation
VALUES (0), (1) tab(col)
|> AGGREGATE COUNT(col) AS count;
+-----+
|count|
+-----+
| 2|
+-----+
-- Aggregation with grouping
VALUES (0, 1), (0, 2) tab(col1, col2)
|> AGGREGATE COUNT(col2) AS count GROUP BY col1;
+----+-----+
|col1|count|
+----+-----+
| 0| 2|
+----+-----+
JOIN
|> [LEFT | RIGHT | FULL | CROSS | SEMI | ANTI | NATURAL | LATERAL] JOIN <table> [ON <condition> | USING(col, ...)]
두 입력의 행을 결합해, 파이프 입력 테이블과 JOIN 키워드 다음의 테이블 표현식의 필터링된 교차곱을 반환해요. 이는 일반 SQL의 JOIN 절과 유사하게 동작하는데, 파이프 연산자 입력 테이블이 조인의 왼쪽이 되고 테이블 인자가 조인의 오른쪽이 돼요.
LEFT, RIGHT, FULL 같은 표준 조인 수정자는 JOIN 키워드 앞에 지원돼요.
조인 조건(predicate)은 조인의 두 입력 모두의 컬럼을 참조해야 할 수 있어요. 이 경우 두 입력에 같은 이름의 컬럼이 있을 때 컬럼을 구분하기 위해 테이블 별칭을 사용해야 할 수 있어요. AS 연산자는 조인의 왼쪽이 되는 파이프 입력 테이블에 새 별칭을 도입하는 데 유용할 수 있어요. 필요하다면 조인의 오른쪽이 되는 테이블 인자에 별칭을 할당하는 데는 표준 구문을 사용하세요.
예를 들어:
SELECT 0 AS a, 1 AS b
|> AS lhs
|> JOIN VALUES (0, 2) rhs(a, b) ON (lhs.a = rhs.a);
+---+---+---+---+
| a| b| c| d|
+---+---+---+---+
| 0| 1| 0| 2|
+---+---+---+---+
VALUES ('apples', 3), ('bananas', 4) t(item, sales)
|> AS produce_sales
|> LEFT JOIN
(SELECT "apples" AS item, 123 AS id) AS produce_data
USING (item)
|> SELECT produce_sales.item, sales, id;
/*---------+-------+------+
| item | sales | id |
+---------+-------+------+
| apples | 3 | 123 |
| bananas | 4 | NULL |
+---------+-------+------*/
ORDER BY
|> ORDER BY <expr> [ASC | DESC], ...
지시된 대로 정렬한 후 입력 행을 반환해요. NULLS FIRST/LAST를 포함한 표준 수정자가 지원돼요.
예를 들어:
VALUES (0), (1) tab(col)
|> ORDER BY col DESC;
+---+
|col|
+---+
| 1|
| 0|
+---+
UNION, INTERSECT, EXCEPT
|> {UNION | INTERSECT | EXCEPT} {ALL | DISTINCT} (<query>)
입력 테이블 또는 하위 쿼리의 결합된 행에 대해 union 또는 기타 집합 연산을 수행해요.
예를 들어:
VALUES (0), (1) tab(a, b)
|> UNION ALL VALUES (2), (3) tab(c, d);
+---+----+
| a| b|
+---+----+
| 0| 1|
| 2| 3|
+---+----+
TABLESAMPLE
|> TABLESAMPLE <method>(<size> {ROWS | PERCENT})
제공된 샘플링 알고리즘이 선택한 행의 부분집합을 반환해요.
예를 들어:
VALUES (0), (0), (0), (0) tab(col)
|> TABLESAMPLE (1 ROWS);
+---+
|col|
+---+
| 0|
+---+
VALUES (0), (0) tab(col)
|> TABLESAMPLE (100 PERCENT);
+---+
|col|
+---+
| 0|
| 0|
+---+
PIVOT
|> PIVOT (agg_expr FOR col IN (val1, ...))
입력 행이 컬럼이 되도록 피벗된 새 테이블을 반환해요.
예를 들어:
VALUES
("dotNET", 2012, 10000),
("Java", 2012, 20000),
("dotNET", 2012, 5000),
("dotNET", 2013, 48000),
("Java", 2013, 30000)
courseSales(course, year, earnings)
|> PIVOT (
SUM(earnings)
FOR COURSE IN ('dotNET', 'Java')
)
+----+------+------+
|year|dotNET| Java|
+----+------+------+
|2012| 15000| 20000|
|2013| 48000| 30000|
+----+------+------+
UNPIVOT
|> UNPIVOT (value_col FOR key_col IN (col1, ...))
입력 컬럼이 행이 되도록 피벗된 새 테이블을 반환해요.
예를 들어:
VALUES
("dotNET", 2012, 10000),
("Java", 2012, 20000),
("dotNET", 2012, 5000),
("dotNET", 2013, 48000),
("Java", 2013, 30000)
courseSales(course, year, earnings)
|> UNPIVOT (
earningsYear FOR `year` IN (`2012`, `2013`, `2014`)
+--------+------+--------+
| course| year|earnings|
+--------+------+--------+
| Java| 2012| 20000|
| Java| 2013| 30000|
| dotNET| 2012| 15000|
| dotNET| 2013| 48000|
| dotNET| 2014| 22500|
+--------+------+--------+
더 알아보기 (Learn more)
- 아파치 스파크 SQL 파이프 구문 (원문)
- SQL 참조 (원문) — Spark SQL 참조
- Spark SQL 시작하기 (원문) — Spark SQL 가이드