패턴 인식
패턴 인식 (Pattern Recognition)
이벤트 패턴 집합을 검색하는 것은 특히 데이터 스트림에서 흔한 사용 사례입니다. Flink는 이벤트 스트림에서 패턴 감지를 허용하는 복합 이벤트 처리(CEP) 라이브러리를 제공합니다. 또한 Flink의 SQL API는 바로 사용할 수 있는 많은 내장 함수와 규칙 기반 최적화를 가진 쿼리를 표현하는 관계형 방식을 제공합니다.
2016년 12월, 국제 표준화 기구(ISO)는 SQL의 Row Pattern Recognition(ISO/IEC TR 19075-5:2016)을 포함하는 SQL 표준의 새 버전을 출시했습니다. 이를 통해 Flink는 SQL에서 복합 이벤트 처리를 위한 MATCH_RECOGNIZE 절을 사용하여 CEP와 SQL API를 통합할 수 있습니다.
MATCH_RECOGNIZE 절은 다음 작업을 가능하게 합니다.
PARTITION BY와ORDER BY절로 사용되는 데이터를 논리적으로 파티셔닝하고 정렬합니다.PATTERN절로 검색할 행 패턴을 정의합니다. 이 패턴은 정규식과 유사한 구문을 사용합니다.- 행 패턴 변수의 논리적 구성 요소는
DEFINE절에 지정됩니다. - SQL 쿼리의 다른 부분에서 사용할 수 있는 표현식인 measure를
MEASURES절에서 정의합니다.
다음 예시는 기본 패턴 인식의 구문을 보여줍니다.
SELECT T.aid, T.bid, T.cid
FROM MyTable
MATCH_RECOGNIZE (
PARTITION BY userid
ORDER BY proctime
MEASURES
A.id AS aid,
B.id AS bid,
C.id AS cid
PATTERN (A B C)
DEFINE
A AS name = 'a',
B AS name = 'b',
C AS name = 'c'
) AS T
이 페이지는 각 키워드를 더 자세히 설명하고 더 복잡한 예시를 보여줍니다.
Flink의
MATCH_RECOGNIZE절 구현은 전체 표준의 일부입니다. 다음 섹션에 문서화된 기능만 지원됩니다. 커뮤니티 피드백에 따라 추가 기능이 지원될 수 있으며, 알려진 제한 사항도 확인해 주세요.
출처: 문서
본문
소개 및 예시
설치 가이드
패턴 인식 기능은 내부적으로 Apache Flink의 CEP 라이브러리를 사용합니다. MATCH_RECOGNIZE 절을 사용하려면 Maven 프로젝트에 라이브러리를 의존성으로 추가해야 합니다.
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-cep</artifactId>
<version>2.3.0</version>
</dependency>
또한 의존성을 클러스터 클래스패스에 추가할 수도 있습니다(자세한 내용은 dependency section 참조).
SQL Client에서 MATCH_RECOGNIZE 절을 사용하려면 모든 의존성이 기본적으로 포함되므로 아무것도 할 필요가 없습니다.
SQL 의미론
모든 MATCH_RECOGNIZE 쿼리는 다음 절로 구성됩니다.
- PARTITION BY - 테이블의 논리적 파티셔닝을 정의합니다.
GROUP BY연산과 유사합니다. - ORDER BY - 들어오는 행이 어떻게 정렬되어야 하는지 지정합니다. 패턴은 순서에 의존하므로 필수적입니다.
- MEASURES - 절의 출력을 정의합니다.
SELECT절과 유사합니다. - ONE ROW PER MATCH - 매치당 몇 개의 행을 생성할지 정의하는 출력 모드입니다.
- AFTER MATCH SKIP - 다음 매치가 시작될 위치를 지정합니다. 단일 이벤트가 몇 개의 서로 다른 매치에 속할 수 있는지 제어하는 방법이기도 합니다.
- PATTERN - 정규식과 유사한 구문으로 검색할 패턴을 구성할 수 있게 합니다.
- DEFINE - 패턴 변수가 충족해야 하는 조건을 정의하는 섹션입니다.
주의: 현재 MATCH_RECOGNIZE 절은 append table에만 적용할 수 있습니다. 또한 항상 append table을 생성합니다.
예시
예시에서 Ticker 테이블이 등록되었다고 가정합니다. 이 테이블은 특정 시점의 주식 가격을 포함합니다.
테이블은 다음 스키마를 가집니다.
Ticker
|-- symbol: String # symbol of the stock
|-- price: Long # price of the stock
|-- tax: Long # tax liability of the stock
|-- rowtime: TimeIndicatorTypeInfo(rowtime) # point in time when the change to those values happened
단순화를 위해 단일 주식 ACME의 들어오는 데이터만 고려합니다. 티커는 행이 지속적으로 추가되는 다음 표와 유사할 수 있습니다.
symbol rowtime price tax
====== ==================== ======= =======
'ACME' '01-Apr-11 10:00:00' 12 1
'ACME' '01-Apr-11 10:00:01' 17 2
'ACME' '01-Apr-11 10:00:02' 19 1
'ACME' '01-Apr-11 10:00:03' 21 3
'ACME' '01-Apr-11 10:00:04' 25 2
'ACME' '01-Apr-11 10:00:05' 18 1
'ACME' '01-Apr-11 10:00:06' 15 1
'ACME' '01-Apr-11 10:00:07' 14 2
'ACME' '01-Apr-11 10:00:08' 24 2
'ACME' '01-Apr-11 10:00:09' 25 2
'ACME' '01-Apr-11 10:00:10' 19 1
이제 작업은 단일 티커의 지속적으로 감소하는 가격 기간을 찾는 것입니다. 이를 위해 다음과 같은 쿼리를 작성할 수 있습니다.
SELECT *
FROM Ticker
MATCH_RECOGNIZE (
PARTITION BY symbol
ORDER BY rowtime
MEASURES
START_ROW.rowtime AS start_tstamp,
LAST(PRICE_DOWN.rowtime) AS bottom_tstamp,
LAST(PRICE_UP.rowtime) AS end_tstamp
ONE ROW PER MATCH
AFTER MATCH SKIP TO LAST PRICE_UP
PATTERN (START_ROW PRICE_DOWN+ PRICE_UP)
DEFINE
PRICE_DOWN AS
(LAST(PRICE_DOWN.price, 1) IS NULL AND PRICE_DOWN.price < START_ROW.price) OR
PRICE_DOWN.price < LAST(PRICE_DOWN.price, 1),
PRICE_UP AS
PRICE_UP.price > LAST(PRICE_DOWN.price, 1)
) MR;
쿼리는 Ticker 테이블을 symbol 열로 파티셔닝하고 rowtime 시간 속성으로 정렬합니다.
PATTERN 절은 시작 이벤트 START_ROW가 하나 이상의 PRICE_DOWN 이벤트 뒤에 오고 PRICE_UP 이벤트로 끝나는 패턴에 관심이 있음을 지정합니다. 그러한 패턴을 찾을 수 있다면 AFTER MATCH SKIP TO LAST 절이 나타내는 대로 다음 패턴 매치가 마지막 PRICE_UP 이벤트에서 탐색됩니다.
DEFINE 절은 PRICE_DOWN과 PRICE_UP 이벤트에 대해 충족해야 하는 조건을 지정합니다. START_ROW 패턴 변수는 존재하지 않지만 항상 TRUE로 평가되는 암시적 조건을 가집니다.
패턴 변수 PRICE_DOWN은 PRICE_DOWN 조건을 충족한 마지막 행보다 작은 가격을 가진 행으로 정의됩니다. 초기 경우 또는 PRICE_DOWN 조건을 충족한 마지막 행이 없을 때는 행의 가격이 패턴의 앞선 행(START_ROW로 참조)의 가격보다 작아야 합니다.
패턴 변수 PRICE_UP은 PRICE_DOWN 조건을 충족한 마지막 행보다 큰 가격을 가진 행으로 정의됩니다.
이 쿼리는 주식 가격이 지속적으로 감소한 각 기간에 대한 요약 행을 생성합니다.
출력 행의 정확한 표현은 쿼리의 MEASURES 부분에서 정의됩니다. 출력 행 수는 ONE ROW PER MATCH 출력 모드로 정의됩니다.
symbol start_tstamp bottom_tstamp end_tstamp
========= ================== ================== ==================
ACME 01-APR-11 10:00:04 01-APR-11 10:00:07 01-APR-11 10:00:08
결과 행은 01-APR-11 10:00:04에 시작해 01-APR-11 10:00:07에 최저 가격을 달성하고 01-APR-11 10:00:08에 다시 상승한 가격 하락 기간을 설명합니다.
파티셔닝 (Partitioning)
파티셔닝된 데이터에서 패턴을 찾는 것이 가능합니다. 예: 단일 티커나 특정 사용자의 추세. 이는 PARTITION BY 절로 표현할 수 있습니다. 이 절은 집계에 GROUP BY를 사용하는 것과 유사합니다.
들어오는 데이터를 파티셔닝하는 것이 매우 권장됩니다. 그렇지 않으면 MATCH_RECOGNIZE 절이 전역 정렬을 보장하기 위해 비병렬 연산자로 변환되기 때문입니다.
이벤트 순서 (Order of Events)
Apache Flink는 처리 시간 또는 이벤트 시간 기반으로 패턴 검색을 허용합니다.
이벤트 시간의 경우 이벤트는 내부 패턴 상태 머신에 전달되기 전에 정렬됩니다. 결과적으로 행이 테이블에 추가되는 순서와 관계없이 생성된 출력은 올바릅니다. 대신 패턴은 각 행에 포함된 시간이 지정한 순서로 평가됩니다.
MATCH_RECOGNIZE 절은 ORDER BY 절의 첫 번째 인수로 오름차순 정렬된 시간 속성을 가정합니다.
예시 Ticker 테이블의 경우 ORDER BY rowtime ASC, price DESC 같은 정의는 유효하지만 ORDER BY price, rowtime 또는 ORDER BY rowtime DESC, price ASC는 유효하지 않습니다.
Define & Measures
DEFINE과 MEASURES 키워드는 간단한 SQL 쿼리의 WHERE와 SELECT 절과 유사한 의미를 가집니다.
MEASURES 절은 일치하는 패턴의 출력에 무엇이 포함될지 정의합니다. 열을 투영하고 평가할 표현식을 정의할 수 있습니다. 생성되는 행 수는 출력 모드 설정에 따라 달라집니다.
DEFINE 절은 행이 해당 패턴 변수로 분류되기 위해 충족해야 하는 조건을 지정합니다. 패턴 변수에 대해 조건이 정의되지 않으면 모든 행에 대해 true로 평가되는 기본 조건이 사용됩니다.
이 절에서 사용할 수 있는 표현식에 대한 더 자세한 설명은 이벤트 스트림 탐색 섹션을 참조하세요.
집계 (Aggregations)
집계는 DEFINE과 MEASURES 절에서 사용할 수 있습니다. 내장 및 사용자 정의 user defined 함수가 모두 지원됩니다.
집계 함수는 매치에 매핑된 각 행 하위 집합에 적용됩니다. 이러한 하위 집합이 어떻게 평가되는지 이해하려면 이벤트 스트림 탐색 섹션을 참조하세요.
다음 예시의 작업은 티커의 평균 가격이 특정 임계값 아래로 내려가지 않은 가장 긴 기간을 찾는 것입니다. MATCH_RECOGNIZE가 집계로 얼마나 표현력 있게 될 수 있는지 보여줍니다. 이 작업은 다음 쿼리로 수행할 수 있습니다.
SELECT *
FROM Ticker
MATCH_RECOGNIZE (
PARTITION BY symbol
ORDER BY rowtime
MEASURES
FIRST(A.rowtime) AS start_tstamp,
LAST(A.rowtime) AS end_tstamp,
AVG(A.price) AS avgPrice
ONE ROW PER MATCH
AFTER MATCH SKIP PAST LAST ROW
PATTERN (A+ B)
DEFINE
A AS AVG(A.price) < 15
) MR;
이 쿼리와 다음 입력 값이 주어지면:
symbol rowtime price tax
====== ==================== ======= =======
'ACME' '01-Apr-11 10:00:00' 12 1
'ACME' '01-Apr-11 10:00:01' 17 2
'ACME' '01-Apr-11 10:00:02' 13 1
'ACME' '01-Apr-11 10:00:03' 16 3
'ACME' '01-Apr-11 10:00:04' 25 2
'ACME' '01-Apr-11 10:00:05' 2 1
'ACME' '01-Apr-11 10:00:06' 4 1
'ACME' '01-Apr-11 10:00:07' 10 2
'ACME' '01-Apr-11 10:00:08' 15 2
'ACME' '01-Apr-11 10:00:09' 25 2
'ACME' '01-Apr-11 10:00:10' 25 1
'ACME' '01-Apr-11 10:00:11' 30 1
쿼리는 평균 가격이 15를 초과하지 않는 한 이벤트를 패턴 변수 A의 일부로 축적합니다. 예를 들어 그러한 한계 초과는 01-Apr-11 10:00:04에 발생합니다. 다음 기간은 01-Apr-11 10:00:11에 평균 가격 15를 다시 초과합니다. 따라서 해당 쿼리의 결과는 다음과 같습니다.
symbol start_tstamp end_tstamp avgPrice
========= ================== ================== ============
ACME 01-APR-11 10:00:00 01-APR-11 10:00:03 14.5
ACME 01-APR-11 10:00:05 01-APR-11 10:00:10 13.5
집계는 단일 패턴 변수를 참조하는 경우에만 표현식에 적용할 수 있습니다. 따라서 SUM(A.price * A.tax)는 유효하지만 AVG(A.price * B.tax)는 유효하지 않습니다.
DISTINCT집계는 지원되지 않습니다.
패턴 정의 (Defining a Pattern)
MATCH_RECOGNIZE 절은 널리 쓰이는 정규식 구문과 다소 유사한 강력하고 표현력 있는 구문으로 이벤트 스트림에서 패턴을 검색할 수 있게 합니다.
모든 패턴은 패턴 변수라고 하는 기본 구성 요소로 구성되며, 여기에 연산자(수량자 및 기타 수정자)를 적용할 수 있습니다. 전체 패턴은 괄호로 묶어야 합니다.
예시 패턴은 다음과 같을 수 있습니다.
PATTERN (A B+ C* D)
다음 연산자를 사용할 수 있습니다.
- 연결(Concatenation) -
(A B)같은 패턴은A와B사이의 연속성이 엄격함을 의미합니다. 따라서 사이에A나B에 매핑되지 않은 행이 있을 수 없습니다. - 수량자(Quantifiers) - 패턴 변수에 매핑될 수 있는 행 수를 수정합니다.
*— 0개 이상의 행+— 1개 이상의 행?— 0개 또는 1개의 행{ n }— 정확히 n개의 행 (n > 0){ n, }— n개 이상의 행 (n ≥ 0){ n, m }— n과 m(포함) 사이의 행 (0 ≤ n ≤ m, 0 < m){ , m }— 0과 m(포함) 사이의 행 (m > 0)
잠재적으로 빈 매치를 생성할 수 있는 패턴은 지원되지 않습니다. 그러한 패턴의 예는
PATTERN (A*),PATTERN (A? B*),PATTERN (A{0,} B{0,} C*)등입니다.
Greedy 및 Reluctant 수량자
각 수량자는 greedy(기본 동작) 또는 reluctant일 수 있습니다. Greedy 수량자는 가능한 많은 행을 매칭하려고 시도하고 reluctant 수량자는 가능한 적게 매칭하려고 시도합니다.
차이를 설명하기 위해 B 변수에 greedy 수량자가 적용된 다음 예시 쿼리를 볼 수 있습니다.
SELECT *
FROM Ticker
MATCH_RECOGNIZE(
PARTITION BY symbol
ORDER BY rowtime
MEASURES
C.price AS lastPrice
ONE ROW PER MATCH
AFTER MATCH SKIP PAST LAST ROW
PATTERN (A B* C)
DEFINE
A AS A.price > 10,
B AS B.price < 15,
C AS C.price > 12
)
다음 입력이 있다면:
symbol tax price rowtime
======= ===== ======== =====================
XYZ 1 10 2018-09-17 10:00:02
XYZ 2 11 2018-09-17 10:00:03
XYZ 1 12 2018-09-17 10:00:04
XYZ 2 13 2018-09-17 10:00:05
XYZ 1 14 2018-09-17 10:00:06
XYZ 2 16 2018-09-17 10:00:07
위 패턴은 다음 출력을 생성합니다.
symbol lastPrice
======== ===========
XYZ 16
B*가 B*?로 수정된(즉 B*가 reluctant여야 함) 같은 쿼리는 다음을 생성합니다.
symbol lastPrice
======== ===========
XYZ 13
XYZ 16
패턴 변수 B는 가격 12, 13, 14의 행을 삼키는 대신 가격 12의 행에만 매칭됩니다.
패턴의 마지막 변수에 greedy 수량자를 사용하는 것은 불가능합니다. 따라서 (A B*) 같은 패턴은 허용되지 않습니다. 이는 B의 부정 조건을 가진 인공 상태(예: C)를 도입하여 쉽게 해결할 수 있습니다. 따라서 다음 같은 쿼리를 사용할 수 있습니다.
PATTERN (A B* C)
DEFINE
A AS condA(),
B AS condB(),
C AS NOT condB()
주의: 선택적 reluctant 수량자(A?? 또는 A{0,1}?)는 현재 지원되지 않습니다.
시간 제약 (Time constraint)
특히 스트리밍 사용 사례에서 패턴이 주어진 기간 내에 완료되는 것이 종종 필요합니다. 이는 greedy 수량자의 경우에도 Flink가 내부적으로 유지해야 하는 전체 상태 크기를 제한할 수 있게 합니다.
따라서 Flink SQL은 패턴의 시간 제약을 정의하는 추가적인(비표준 SQL) WITHIN 절을 지원합니다. 이 절은 PATTERN 절 뒤에 정의할 수 있으며 밀리초 해상도의 간격을 받습니다.
잠재적 매치의 첫 번째와 마지막 이벤트 사이의 시간이 주어진 값보다 길면 그러한 매치는 결과 테이블에 추가되지 않습니다.
참고: WITHIN 절은 Flink의 효율적인 메모리 관리에 도움이 되므로 일반적으로 사용이 권장됩니다. 임계값에 도달하면 기본 상태를 정리할 수 있습니다.
주의: 그러나 WITHIN 절은 SQL 표준의 일부가 아닙니다. 시간 제약을 처리하는 권장 방법은 향후 변경될 수 있습니다.
WITHIN 절의 사용은 다음 예시 쿼리에서 설명됩니다.
SELECT *
FROM Ticker
MATCH_RECOGNIZE(
PARTITION BY symbol
ORDER BY rowtime
MEASURES
C.rowtime AS dropTime,
A.price - C.price AS dropDiff
ONE ROW PER MATCH
AFTER MATCH SKIP PAST LAST ROW
PATTERN (A B* C) WITHIN INTERVAL '1' HOUR
DEFINE
B AS B.price > A.price - 10,
C AS C.price < A.price - 10
)
쿼리는 1시간 간격 내에 발생하는 10의 가격 하락을 감지합니다.
쿼리가 다음 티커 데이터를 분석하는 데 사용된다고 가정합니다.
symbol rowtime price tax
====== ==================== ======= =======
'ACME' '01-Apr-11 10:00:00' 20 1
'ACME' '01-Apr-11 10:20:00' 17 2
'ACME' '01-Apr-11 10:40:00' 18 1
'ACME' '01-Apr-11 11:00:00' 11 3
'ACME' '01-Apr-11 11:20:00' 14 2
'ACME' '01-Apr-11 11:40:00' 9 1
'ACME' '01-Apr-11 12:00:00' 15 1
'ACME' '01-Apr-11 12:20:00' 14 2
'ACME' '01-Apr-11 12:40:00' 24 2
'ACME' '01-Apr-11 13:00:00' 1 2
'ACME' '01-Apr-11 13:20:00' 19 1
쿼리는 다음 결과를 생성합니다.
symbol dropTime dropDiff
====== ==================== =============
'ACME' '01-Apr-11 13:00:00' 14
결과 행은 15(01-Apr-11 12:00:00)에서 1(01-Apr-11 13:00:00)로의 가격 하락을 나타냅니다. dropDiff 열은 가격 차이를 포함합니다.
예를 들어 11(01-Apr-11 10:00:00와 01-Apr-11 11:40:00 사이)로 가격이 더 크게 하락하더라도, 두 이벤트 사이의 시간 차이가 1시간보다 크므로 매치를 생성하지 않습니다.
출력 모드 (Output Mode)
출력 모드는 찾은 모든 매치에 대해 몇 개의 행을 방출해야 하는지 설명합니다. SQL 표준은 두 가지 모드를 설명합니다.
ALL ROWS PER MATCHONE ROW PER MATCH.
현재 지원되는 유일한 출력 모드는 ONE ROW PER MATCH이며, 이는 찾은 각 매치에 대해 항상 하나의 출력 요약 행을 생성합니다.
출력 행의 스키마는 [파티셔닝 열] + [measure 열]의 연결이 해당 순서대로 됩니다.
다음 예시는 다음과 같이 정의된 쿼리의 출력을 보여줍니다.
SELECT *
FROM Ticker
MATCH_RECOGNIZE(
PARTITION BY symbol
ORDER BY rowtime
MEASURES
FIRST(A.price) AS startPrice,
LAST(A.price) AS topPrice,
B.price AS lastPrice
ONE ROW PER MATCH
PATTERN (A+ B)
DEFINE
A AS LAST(A.price, 1) IS NULL OR A.price > LAST(A.price, 1),
B AS B.price < LAST(A.price)
)
다음 입력 행에 대해:
symbol tax price rowtime
======== ===== ======== =====================
XYZ 1 10 2018-09-17 10:00:02
XYZ 2 12 2018-09-17 10:00:03
XYZ 1 13 2018-09-17 10:00:04
XYZ 2 11 2018-09-17 10:00:05
쿼리는 다음 출력을 생성합니다.
symbol startPrice topPrice lastPrice
======== ============ ========== ===========
XYZ 10 13 11
패턴 인식은 symbol 열로 파티셔닝됩니다. MEASURES 절에서 명시적으로 언급되지 않았지만 파티셔닝된 열이 결과의 시작에 추가됩니다.
패턴 탐색 (Pattern Navigation)
DEFINE과 MEASURES 절은 (잠재적으로) 패턴과 일치하는 행 목록 내에서 탐색을 허용합니다.
이 섹션은 조건 선언 또는 출력 결과 생성을 위한 이 탐색을 논의합니다.
패턴 변수 참조
패턴 변수 참조는 DEFINE 또는 MEASURES 절에서 특정 패턴 변수에 매핑된 행 집합을 참조할 수 있게 합니다.
예를 들어 표현식 A.price는 지금까지 A에 매핑된 행 집합에 현재 행을 A에 매칭하려고 하면 현재 행을 더한 집합을 설명합니다. DEFINE/MEASURES 절의 표현식이 단일 행을 요구하면(예: A.price 또는 A.price > 10) 해당 집합에 속하는 마지막 값을 선택합니다.
패턴 변수가 지정되지 않으면(예: SUM(price)) 표현식은 패턴의 모든 변수를 참조하는 기본 패턴 변수 *를 참조합니다. 즉 지금까지 어떤 변수에도 매핑된 모든 행에 현재 행을 더한 목록을 만듭니다.
예시
더 철저한 예시를 위해 다음 패턴과 해당 조건을 살펴볼 수 있습니다.
PATTERN (A B+)
DEFINE
A AS A.price >= 10,
B AS B.price > A.price AND SUM(price) < 100 AND SUM(B.price) < 80
다음 표는 각 들어오는 이벤트에 대해 이러한 조건이 어떻게 평가되는지 설명합니다.
표는 다음 열로 구성됩니다.
#- 목록[A.price]/[B.price]/[price]에서 들어오는 행을 고유하게 식별하는 행 식별자입니다.price- 들어오는 행의 가격입니다.[A.price]/[B.price]/[price]-DEFINE절에서 조건을 평가하는 데 사용되는 행 목록을 설명합니다.Classifier- 행이 매핑된 패턴 변수를 나타내는 현재 행의 분류기입니다.A.price/B.price/SUM(price)/SUM(B.price)- 해당 표현식이 평가된 후의 결과를 설명합니다.
| # | price | Classifier | [A.price] | [B.price] | [price] | A.price | B.price | SUM(price) | SUM(B.price) |
|---|---|---|---|---|---|---|---|---|---|
| #1 | 10 | -> A | #1 | - | - | 10 | - | - | - |
| #2 | 15 | -> B | #1 | #2 | #1, #2 | 10 | 15 | 25 | 15 |
| #3 | 20 | -> B | #1 | #2, #3 | #1, #2, #3 | 10 | 20 | 45 | 35 |
| #4 | 31 | -> B | #1 | #2, #3, #4 | #1, #2, #3, #4 | 10 | 31 | 76 | 66 |
| #5 | 35 | #1 | #2, #3, #4, #5 | #1, #2, #3, #4, #5 | 10 | 35 | 111 | 101 |
표에서 볼 수 있듯이 첫 번째 행은 패턴 변수 A에 매핑되고 후속 행은 패턴 변수 B에 매핑됩니다. 그러나 마지막 행은 매핑된 모든 행의 합 SUM(price)과 B의 모든 행의 합이 지정된 임계값을 초과하므로 B 조건을 충족하지 못합니다.
논리적 오프셋 (Logical Offsets)
논리적 오프셋은 특정 패턴 변수에 매핑된 이벤트 내에서 탐색을 가능하게 합니다. 이는 두 가지 해당 함수로 표현할 수 있습니다.
| Offset functions | Description |
|---|---|
LAST(variable.field, n) |
변수의 n번째 마지막 요소에 매핑된 이벤트의 필드 값을 반환합니다. 카운트는 매핑된 마지막 요소에서 시작합니다. |
FIRST(variable.field, n) |
변수의 n번째 요소에 매핑된 이벤트의 필드 값을 반환합니다. 카운트는 매핑된 첫 번째 요소에서 시작합니다. |
예시
더 철저한 예시를 위해 다음 패턴과 해당 조건을 살펴볼 수 있습니다.
PATTERN (A B+)
DEFINE
A AS A.price >= 10,
B AS (LAST(B.price, 1) IS NULL OR B.price > LAST(B.price, 1)) AND
(LAST(B.price, 2) IS NULL OR B.price > 2 * LAST(B.price, 2))
다음 표는 각 들어오는 이벤트에 대해 이러한 조건이 어떻게 평가되는지 설명합니다.
표는 다음 열로 구성됩니다.
price- 들어오는 행의 가격입니다.Classifier- 행이 매핑된 패턴 변수를 나타내는 현재 행의 분류기입니다.LAST(B.price, 1)/LAST(B.price, 2)- 해당 표현식이 평가된 후의 결과를 설명합니다.
| price | Classifier | LAST(B.price, 1) | LAST(B.price, 2) | Comment |
|---|---|---|---|---|
| 10 | -> A | |||
| 15 | -> B | null | null | LAST(B.price, 1)이 B에 아직 아무것도 매핑되지 않았으므로 null입니다. |
| 20 | -> B | 15 | null | |
| 31 | -> B | 20 | 15 | |
| 35 | 31 | 20 | 35 < 2 * 20이므로 매핑되지 않음. |
논리적 오프셋과 함께 기본 패턴 변수를 사용하는 것도 의미가 있을 수 있습니다.
이 경우 오프셋은 지금까지 매핑된 모든 행을 고려합니다.
PATTERN (A B? C)
DEFINE
B AS B.price < 20,
C AS LAST(price, 1) < C.price
| price | Classifier | LAST(price, 1) | Comment |
|---|---|---|---|
| 10 | -> A | ||
| 15 | -> B | ||
| 20 | -> C | 15 | LAST(price, 1)은 B 변수에 매핑된 행의 가격으로 평가됩니다. |
두 번째 행이 B 변수에 매핑되지 않았다면 다음 결과가 있습니다.
| price | Classifier | LAST(price, 1) | Comment |
|---|---|---|---|
| 10 | -> A | ||
| 20 | -> C | 10 | LAST(price, 1)은 A 변수에 매핑된 행의 가격으로 평가됩니다. |
FIRST/LAST 함수의 첫 번째 인자에서 여러 패턴 변수 참조를 사용하는 것도 가능합니다. 이렇게 하면 여러 열에 접근하는 표현식을 작성할 수 있습니다. 그러나 모두 같은 패턴 변수를 사용해야 합니다. 즉 LAST/FIRST 함수의 값은 단일 행에서 계산되어야 합니다.
따라서 LAST(A.price * A.tax)는 사용할 수 있지만 LAST(A.price * B.tax) 같은 표현식은 허용되지 않습니다.
After Match 전략
AFTER MATCH SKIP 절은 완전한 매치를 찾은 후 새로운 매칭 절차를 시작할 위치를 지정합니다.
네 가지 전략이 있습니다.
SKIP PAST LAST ROW- 현재 매치의 마지막 행 다음 행에서 패턴 매칭을 재개합니다.SKIP TO NEXT ROW- 매치의 시작 행 다음 행부터 새 매치를 계속 검색합니다.SKIP TO LAST variable- 지정된 패턴 변수에 매핑된 마지막 행에서 패턴 매칭을 재개합니다.SKIP TO FIRST variable- 지정된 패턴 변수에 매핑된 첫 번째 행에서 패턴 매칭을 재개합니다.
이것은 또한 단일 이벤트가 몇 개의 매치에 속할 수 있는지 지정하는 방법입니다. 예를 들어 SKIP PAST LAST ROW 전략을 사용하면 모든 이벤트는 최대 하나의 매치에 속할 수 있습니다.
예시
이 전략들의 차이를 더 잘 이해하기 위해 다음 예시를 살펴볼 수 있습니다.
다음 입력 행에 대해:
symbol tax price rowtime
======== ===== ======= =====================
XYZ 1 7 2018-09-17 10:00:01
XYZ 2 9 2018-09-17 10:00:02
XYZ 1 10 2018-09-17 10:00:03
XYZ 2 5 2018-09-17 10:00:04
XYZ 2 10 2018-09-17 10:00:05
XYZ 2 7 2018-09-17 10:00:06
XYZ 2 14 2018-09-17 10:00:07
다양한 전략으로 다음 쿼리를 평가합니다.
SELECT *
FROM Ticker
MATCH_RECOGNIZE(
PARTITION BY symbol
ORDER BY rowtime
MEASURES
SUM(A.price) AS sumPrice,
FIRST(rowtime) AS startTime,
LAST(rowtime) AS endTime
ONE ROW PER MATCH
[AFTER MATCH STRATEGY]
PATTERN (A+ C)
DEFINE
A AS SUM(A.price) < 30
)
쿼리는 A에 매핑된 모든 행의 가격 합과 전체 매치의 첫 번째 및 마지막 타임스탬프를 반환합니다.
쿼리는 사용된 AFTER MATCH 전략에 따라 다른 결과를 생성합니다.
AFTER MATCH SKIP PAST LAST ROW
symbol sumPrice startTime endTime
======== ========== ===================== =====================
XYZ 26 2018-09-17 10:00:01 2018-09-17 10:00:04
XYZ 17 2018-09-17 10:00:05 2018-09-17 10:00:07
첫 번째 결과는 행 #1, #2, #3, #4에 매칭되었습니다.
두 번째 결과는 행 #5, #6, #7에 매칭되었습니다.
AFTER MATCH SKIP TO NEXT ROW
symbol sumPrice startTime endTime
======== ========== ===================== =====================
XYZ 26 2018-09-17 10:00:01 2018-09-17 10:00:04
XYZ 24 2018-09-17 10:00:02 2018-09-17 10:00:05
XYZ 25 2018-09-17 10:00:03 2018-09-17 10:00:06
XYZ 22 2018-09-17 10:00:04 2018-09-17 10:00:07
XYZ 17 2018-09-17 10:00:05 2018-09-17 10:00:07
다시 첫 번째 결과는 행 #1, #2, #3, #4에 매칭되었습니다.
이전 전략과 비교해 다음 매치는 행 #2를 다시 포함합니다. 따라서 두 번째 결과는 행 #2, #3, #4, #5에 매칭되었습니다.
세 번째 결과는 행 #3, #4, #5, #6에 매칭되었습니다.
네 번째 결과는 행 #4, #5, #6, #7에 매칭되었습니다.
마지막 결과는 행 #5, #6, #7에 매칭되었습니다.
AFTER MATCH SKIP TO LAST A
symbol sumPrice startTime endTime
======== ========== ===================== =====================
XYZ 26 2018-09-17 10:00:01 2018-09-17 10:00:04
XYZ 25 2018-09-17 10:00:03 2018-09-17 10:00:06
XYZ 17 2018-09-17 10:00:05 2018-09-17 10:00:07
다시 첫 번째 결과는 행 #1, #2, #3, #4에 매칭되었습니다.
이전 전략과 비교해 다음 매치는 행 #3(A에 매핑됨)만 다시 포함합니다. 따라서 두 번째 결과는 행 #3, #4, #5, #6에 매칭되었습니다.
마지막 결과는 행 #5, #6, #7에 매칭되었습니다.
AFTER MATCH SKIP TO FIRST A
이 조합은 항상 마지막 매치와 같은 곳에서 새 매치를 시작하려고 하므로 런타임 예외를 생성합니다. 이는 무한 루프를 생성하고 따라서 금지됩니다.
SKIP TO FIRST/LAST variable 전략의 경우 그 변수에 매핑된 행이 없을 수 있음을 명심해야 합니다(예: 패턴 A*). 그러한 경우 표준이 매칭을 계속하려면 유효한 행이 필요하므로 런타임 예외가 던져집니다.
시간 속성
MATCH_RECOGNIZE 위에 일부 후속 쿼리를 적용하려면 시간 속성을 사용해야 할 수 있습니다. 그것을 선택하기 위해 사용할 수 있는 두 가지 함수가 있습니다.
| Function | Description |
|---|---|
MATCH_ROWTIME([rowtime_field]) |
주어진 패턴에 매핑된 마지막 행의 타임스탬프를 반환합니다. 이 함수는 0개 또는 1개의 피연산자를 받으며, 피연산자는 rowtime 속성이 있는 필드 참조입니다. 피연산자가 없으면 TIMESTAMP 타입의 rowtime 속성을 반환합니다. 그렇지 않으면 반환 타입은 피연산자 타입과 같습니다. 결과 속성은 rowtime 속성으로, interval joins과 group window 또는 over window 집계 같은 후속 시간 기반 연산에 사용할 수 있습니다. |
MATCH_PROCTIME() |
interval joins과 group window 또는 over window 집계 같은 후속 시간 기반 연산에 사용할 수 있는 proctime 속성을 반환합니다. |
메모리 소비 제어
MATCH_RECOGNIZE 쿼리를 작성할 때 메모리 소비는 중요한 고려 사항입니다. 잠재적 매치의 공간이 너비 우선(breadth-first) 방식으로 구축되기 때문입니다. 이를 염두에 두고 패턴이 끝날 수 있도록 해야 합니다. 바람직하게는 매치에 매핑된 행이 메모리에 맞아야 하므로 합리적인 수로 말이죠.
예를 들어 패턴은 상한 없이 모든 단일 행을 받아들이는 수량자를 가져서는 안 됩니다. 그러한 패턴은 다음과 같을 수 있습니다.
PATTERN (A B+ C)
DEFINE
A as A.price > 10,
C as C.price > 20
쿼리는 모든 들어오는 행을 B 변수에 매핑하고 따라서 결코 끝나지 않습니다. 이 쿼리는 예를 들어 C의 조건을 부정하여 고칠 수 있습니다.
PATTERN (A B+ C)
DEFINE
A as A.price > 10,
B as B.price <= 20,
C as C.price > 20
또는 reluctant 수량자를 사용합니다.
PATTERN (A B+? C)
DEFINE
A as A.price > 10,
C as C.price > 20
주의: MATCH_RECOGNIZE 절은 구성된 상태 보존 시간(state retention time)을 사용하지 않습니다. 이 목적에는 WITHIN 절을 사용하는 것이 좋습니다.
알려진 제한 사항
Flink의 MATCH_RECOGNIZE 절 구현은 진행 중인 작업이며 SQL 표준의 일부 기능은 아직 지원되지 않습니다.
지원되지 않는 기능은 다음과 같습니다.
- 패턴 표현식:
- 패턴 그룹 - 예를 들어 수량자를 패턴의 하위 시퀀스에 적용할 수 없습니다. 따라서
(A (B C)+)는 유효한 패턴이 아닙니다. - 변형(Alterations) -
E행을 찾기 전에 하위 시퀀스A B또는C D중 하나를 찾아야 한다는 뜻의PATTERN((A B | C D) E)같은 패턴. PERMUTE연산자 - 적용된 모든 변수의 순열과 동일합니다. 예:PATTERN (PERMUTE (A, B, C))=PATTERN (A B C | A C B | B A C | B C A | C A B | C B A).- 앵커(Anchors) - 파티션의 시작/끝을 나타내는
^, $. 스트리밍 맥락에서는 의미가 없으며 지원되지 않습니다. - 제외(Exclusion) -
PATTERN ({- A -} B)는A가 검색되지만 출력에 참여하지 않음을 의미합니다. 이는ALL ROWS PER MATCH모드에서만 작동합니다. - Reluctant 선택 수량자 -
PATTERN A??는 greedy 선택 수량자만 지원됩니다.
- 패턴 그룹 - 예를 들어 수량자를 패턴의 하위 시퀀스에 적용할 수 없습니다. 따라서
ALL ROWS PER MATCH출력 모드 - 찾은 매치 생성에 참여한 모든 행에 대해 출력 행을 생성합니다. 이는 다음을 의미하기도 합니다.MEASURES절의 유일하게 지원되는 의미론은FINAL입니다.- 행이 매핑된 패턴 변수를 반환하는
CLASSIFIER함수는 아직 지원되지 않습니다.
SUBSET- 패턴 변수의 논리적 그룹을 만들고DEFINE과MEASURES절에서 그 그룹을 사용하는 것을 허용합니다.- 물리적 오프셋 -
PREV/NEXT. 논리적 오프셋 경우와 달리 오직 패턴 변수에 매핑된 것만이 아니라 본 모든 이벤트를 인덱싱합니다. MATCH_RECOGNIZE는 SQL에서만 지원됩니다. Table API에는 동등한 것이 없습니다.- 집계:
- distinct 집계는 지원되지 않습니다.