최신 값 조회하기
최신 값 조회하기 (Query for latest values)
다른 데이터베이스에서 UPSERT로 처리할 수 있는 사용 사례를 Apache Druid에서 다루는 전략을 설명해 드릴게요. 쿼리 시점에는 LATEST_BY 집계를, 숫자 dimension의 경우 삽입 시점에는 "deltas"를 사용할 수 있어요.
출처: 문서
본문
Update data 튜토리얼은 UPSERT 케이스를 포함해 타임스탬프에 따라 데이터를 업데이트하기 위해 배치 작업을 사용하는 방법을 보여 줘요. 하지만 스트리밍 데이터에서는 업데이트로 처리하는 요구 사항을 LATEST_BY 나 deltas로 충족할 수도 있어요.
사전 준비 (Prerequisites)
이 튜토리얼의 단계를 따라가기 전에, Local quickstart에 설명된 대로 Druid를 다운로드하고 로컬 머신에서 실행 중이어야 해요. Druid 클러스터에 데이터를 로드할 필요는 없어요.
Druid에서 데이터를 쿼리하는 방법에 익숙해야 해요. 아직 안 하셨다면 먼저 Query data 튜토리얼을 진행해 주세요.
LATEST_BY로 업데이트된 값 가져오기 (Use LATEST_BY to retrieve updated values)
때로는 어떤 dimension과 관련해 다른 dimension이나 measure의 최신 값을 읽고 싶을 수 있어요. 트랜잭션 데이터베이스에서는 UPSERT로 dimensions나 measures를 유지할 수 있지만, Druid에서는 수집 중에 모든 업데이트나 변경을 append할 수 있어요. LATEST_BY 함수는 다음 유형의 쿼리로 해당 dimension의 가장 최근 값을 얻을 수 있게 해 줘요:
SELECT dimension,
LATEST_BY(changed_dimension, updated_timestamp)
FROM my_table
GROUP BY 1
이 예시에서 update_timestamp 는 "latest" 값을 평가하는 데 사용할 참조 타임스탬프를 나타내요. 이 값은 __time 이거나 다른 타임스탬프일 수 있어요.
예를 들어 사용자의 총 포인트 수를 기록하는 다음 이벤트 테이블을 생각해 볼게요:
__time |
user_id |
points |
|---|---|---|
| 2024-01-01T01:00:00.000Z | funny_bunny1 | 10 |
| 2024-01-01T01:05:00.000Z | funny_bunny1 | 30 |
| 2024-01-01T02:00:00.000Z | funny_bunny1 | 35 |
| 2024-01-01T02:00:00.000Z | silly_monkey2 | 30 |
| 2024-01-01T02:05:00.000Z | silly_monkey2 | 55 |
| 2024-01-01T03:00:00.000Z | funny_bunny1 | 40 |
샘플 데이터 삽입하기 (Insert sample data)
Druid 웹 콘솔에서 Query 뷰로 이동해 다음 쿼리를 실행해 샘플 데이터를 삽입해 주세요:
REPLACE INTO "latest_by_tutorial1" OVERWRITE ALL
WITH "ext" AS (
SELECT *
FROM TABLE(
EXTERN(
'{"type":"inline","data":"{\"timestamp\":\"2024-01-01T01:00:00Z\",\"user_id\":\"funny_bunny1\", \"points\":10}\n{\"timestamp\":\"2024-01-01T01:05:00Z\",\"user_id\":\"funny_bunny1\", \"points\":30}\n{\"timestamp\": \"2024-01-01T02:00:00Z\",\"user_id\":\"funny_bunny1\", \"points\":35}\n{\"timestamp\":\"2024-01-01T02:00:00Z\",\"user_id\":\"silly_monkey2\", \"points\":30}\n{\"timestamp\":\"2024-01-01T02:05:00Z\",\"user_id\":\"silly_monkey2\", \"points\":55}\n{\"timestamp\":\"2024-01-01T03:00:00Z\",\"user_id\":\"funny_bunny1\", \"points\":40}"}',
'{"type":"json"}'
)
)
EXTEND ("timestamp" VARCHAR, "user_id" VARCHAR, "points" BIGINT)
)
SELECT
TIME_PARSE("timestamp") AS "__time",
"user_id",
"points"
FROM "ext"
PARTITIONED BY DAY
다음 쿼리를 실행해 각 user_id 에 대한 가장 최근 points 값을 가져와 보세요:
SELECT user_id,
LATEST_BY("points", "__time") AS latest_points
FROM latest_by_tutorial1
GROUP BY 1
결과는 다음과 같아요:
user_id |
total_points |
|---|---|
| silly_monkey2 | 55 |
| funny_bunny1 | 40 |
예시에서 값은 매번 증가하지만, 이 방법은 값이 변동하더라도 잘 동작해요.
이 쿼리 형태를 추가 처리를 위한 서브쿼리로 사용할 수도 있어요. 하지만 user_id 에 대한 값이 많으면 쿼리가 비쌀 수 있어요.
더 큰 granularity 시간 프레임 안에서 서로 다른 시점의 최신 값을 추적하고 싶다면, 업데이트 시점을 기록할 추가 타임스탬프가 필요해요. 그래야 Druid가 최신 버전을 추적할 수 있어요. 한 시간 프레임 안에서 여러 사용자의 points가 업데이트되는 다음 데이터를 생각해 보세요. __time 은 시간 granularity인 반면 updated_timestamp 는 분 granularity예요:
__time |
updated_timestamp |
user_id |
points |
|---|---|---|---|
| 2024-01-01T01:00:00.000Z | 2024-01-01T01:00:00.000Z | funny_bunny1 | 10 |
| 2024-01-01T01:00:00.000Z | 2024-01-01T01:05:00.000Z | funny_bunny1 | 30 |
| 2024-01-01T02:00:00.000Z | 2024-01-01T02:00:00.000Z | funny_bunny1 | 35 |
| 2024-01-01T02:00:00.000Z | 2024-01-01T02:00:00.000Z | silly_monkey2 | 30 |
| 2024-01-01T02:00:00.000Z | 2024-01-01T02:05:00.000Z | silly_monkey2 | 55 |
| 2024-01-01T03:00:00.000Z | 2024-01-01T03:00:00.000Z | funny_bunny1 | 40 |
샘플 데이터 삽입하기 (Insert sample data)
Query 뷰에서 새 탭을 열고 다음 쿼리를 실행해 샘플 데이터를 삽입해 주세요:
REPLACE INTO "latest_by_tutorial2" OVERWRITE ALL
WITH "ext" AS (
SELECT *
FROM TABLE(
EXTERN(
'{"type":"inline","data":"{\"timestamp\":\"2024-01-01T01:00:00Z\",\"updated_timestamp\":\"2024-01-01T01:00:00Z\",\"user_id\":\"funny_bunny1\", \"points\":10}\n{\"timestamp\":\"2024-01-01T01:05:00Z\",\"updated_timestamp\":\"2024-01-01T01:05:00Z\",\"user_id\":\"funny_bunny1\", \"points\":30}\n{\"timestamp\": \"2024-01-01T02:00:00Z\",\"updated_timestamp\":\"2024-01-01T02:00:00Z\",\"user_id\":\"funny_bunny1\", \"points\":35}\n{\"timestamp\":\"2024-01-01T02:00:00Z\",\"updated_timestamp\":\"2024-01-01T02:00:00Z\",\"user_id\":\"silly_monkey2\", \"points\":30}\n{\"timestamp\":\"2024-01-01T02:00:00Z\",\"updated_timestamp\":\"2024-01-01T02:05:00Z\",\"user_id\":\"silly_monkey2\", \"points\":55}\n{\"timestamp\":\"2024-01-01T03:00:00Z\",\"updated_timestamp\":\"2024-01-01T03:00:00Z\",\"user_id\":\"funny_bunny1\", \"points\":40}"}',
'{"type":"json"}'
)
)
EXTEND ("timestamp" VARCHAR, "updated_timestamp" VARCHAR, "user_id" VARCHAR, "points" BIGINT)
)
SELECT
TIME_PARSE("timestamp") AS "__time",
"updated_timestamp",
"user_id",
"points"
FROM "ext"
PARTITIONED BY DAY
다음 쿼리를 실행해 각 시간별로 사용자별 최신 points 값을 가져와 보세요:
SELECT FLOOR("__time" TO HOUR) AS "hour_time",
"user_id",
LATEST_BY("points", TIME_PARSE(updated_timestamp)) AS "latest_points_hour"
FROM latest_by_tutorial2
GROUP BY 1,2
결과는 다음과 같아요:
hour_time |
user_id |
latest_points_hour |
|---|---|---|
| 2024-01-01T01:00:00.000Z | funny_bunny1 | 20 |
| 2024-01-01T02:00:00.000Z | funny_bunny1 | 5 |
| 2024-01-01T02:00:00.000Z | silly_monkey2 | 25 |
| 2024-01-01T03:00:00.000Z | funny_bunny1 | 10 |
LATEST_BY 는 집계 함수예요. user_id 같은 dimension과 일치하는 업데이트 행이 많지 않을 때는 매우 효율적이지만, 같은 dimension과 일치하는 모든 행을 스캔해요. 사용자가 게임을 백만 번 하는 것처럼 업데이트가 많은 dimension의 경우, 그리고 업데이트가 시기적절한 순서로 도착하지 않는 경우에는 Druid가 user_id 와 일치하는 모든 행을 처리해 최대 타임스탬프를 가진 행을 찾아 최신 데이터를 제공해요.
예를 들어 업데이트가 데이터의 1~5%를 차지한다면 좋은 쿼리 성능을 얻을 수 있어요. 업데이트가 데이터의 50% 이상을 차지하면 쿼리가 느려질 거예요.
이를 완화하려면 수정된 데이터를 새 데이터소스에 다시 인덱싱하는 주기적인 배치 수집 작업을 설정해서, 최신 값을 미리 계산·저장해 grouping 없이 직접 쿼리함으로써 이런 쿼리 비용을 줄일 수 있어요. 최신 데이터에 대한 관점은 다음 새로고침까지는 최신 상태가 아니라는 점에 주의해 주세요.
또는 수집 시점 집계(ingestion-time aggregation)에서 LATEST_BY 를 사용하고 스트리밍 수집으로 업데이트를 rolled up 데이터소스에 append할 수도 있어요. 시간 청크에 append하면 새 세그먼트가 추가되고 데이터가 완벽하게 roll up되지 않아서, 행이 완전한 rollup이 아닌 부분적일 수 있고 부분적으로 roll up된 행이 여럿 생길 수 있어요. 이 경우에도 rolled up 데이터소스를 올바르게 쿼리하려면 여전히 GROUP BY 쿼리를 사용해야 해요. 자동 컴팩션을 튜닝해서 오래된(stale) 행 수를 크게 줄이고 성능을 개선할 수 있어요.
업데이트된 값에 delta 값과 집계 사용하기 (Use delta values and aggregation for updated values)
이벤트에 최신 총값을 append하는 대신, 각 이벤트와 함께 값의 변화량(change)을 기록하고 평소 쓰는 aggregator를 사용할 수도 있어요. 이 방법은 쿼리에서 한 단계의 집계와 grouping을 피할 수 있게 해 줄 수 있어요.
대부분의 애플리케이션에서는 전처리 없이 이벤트 데이터를 바로 Druid로 보낼 수 있어요. 예를 들어 노출(impression) 횟수를 Druid로 보낼 때, 어제부터의 총 노출 횟수를 보내지 말고 최근 노출 횟수만 보내 주세요. 그러면 쿼리 중에 Druid에서 총합을 집계할 수 있어요. Druid는 많은 행을 더하는 데 최적화되어 있어서, 데이터를 배치하거나 사전 집계하는 데 익숙한 사람에게는 직관에 반할 수도 있어요.
예를 들어 집계 컬럼 y 가 있고 다른 dimension x 로 그룹핑해서 SUM으로 집계하는 데이터소스를 생각해 보세요. y 의 값을 3에서 2로 업데이트하고 싶다면 y 에 -1을 삽입해 주세요. 이렇게 하면 x 로 그룹핑된 어떤 쿼리에서도 SUM(y) 집계가 올바르게 돼요. 이 방법은 상당한 성능 이점을 제공할 수 있지만, 집계가 항상 SUM 이어야 한다는 트레이드오프가 있어요.
다른 경우에는 데이터에 대한 업데이트가 이미 원본에 대한 deltas일 수 있으므로, 업데이트를 append하는 데 필요한 데이터 엔지니어링이 간단할 거예요. 이전 예시와 동일한 성능 영향 완화가 적용돼요: 수집 시점에 rollup을 사용하고 지속적인 자동 컴팩션과 결합하세요.
예를 들어 어떤 기간 동안 사용자가 얻거나 잃은 포인트 수를 기록하는 다음 이벤트 테이블을 생각해 보세요:
__time |
user_id |
delta |
|---|---|---|
| 2024-01-01T01:00:00.000Z | funny_bunny1 | 10 |
| 2024-01-01T01:05:00.000Z | funny_bunny1 | 10 |
| 2024-01-01T02:00:00.000Z | funny_bunny1 | 5 |
| 2024-01-01T02:00:00.000Z | silly_monkey2 | 30 |
| 2024-01-01T02:05:00.000Z | silly_monkey2 | -5 |
| 2024-01-01T03:00:00.000Z | funny_bunny1 | 10 |
샘플 데이터 삽입하기 (Insert sample data)
Query 뷰에서 새 탭을 열고 다음 쿼리를 실행해 샘플 데이터를 삽입해 주세요:
REPLACE INTO "delta_tutorial" OVERWRITE ALL
WITH "ext" AS (
SELECT *
FROM TABLE(
EXTERN(
'{"type":"inline","data":"{\"timestamp\":\"2024-01-01T01:00:00Z\",\"user_id\":\"funny_bunny1\", \"points\":10}\n{\"timestamp\":\"2024-01-01T01:05:00Z\",\"user_id\":\"funny_bunny1\", \"points\":10}\n{\"timestamp\": \"2024-01-01T02:00:00Z\",\"user_id\":\"funny_bunny1\", \"points\":5}\n{\"timestamp\":\"2024-01-01T02:00:00Z\",\"user_id\":\"silly_monkey2\", \"points\":30}\n{\"timestamp\":\"2024-01-01T02:05:00Z\",\"user_id\":\"silly_monkey2\", \"points\":-5}\n{\"timestamp\":\"2024-01-01T03:00:00Z\",\"user_id\":\"funny_bunny1\", \"points\":10}"}',
'{"type":"json"}'
)
)
EXTEND ("timestamp" VARCHAR, "user_id" VARCHAR, "points" BIGINT)
)
SELECT
TIME_PARSE("timestamp") AS "__time",
"user_id",
"points" AS "delta"
FROM "ext"
PARTITIONED BY DAY
다음 쿼리는 두 번째 LATEST_BY 예시와 동일한 시간별 포인트 값을 반환해요:
SELECT FLOOR("__time" TO HOUR) as "hour_time",
"user_id",
SUM("delta") AS "latest_points_hour"
FROM "delta_tutorial"
GROUP BY 1,2
더 알아보기 (Learn more)
자세한 내용은 다음 주제를 참고해 주세요:
- Update data — Druid에서 데이터를 업데이트하는 튜토리얼.
- Data updates — Druid에서 데이터를 업데이트하는 방법 개요.