[アップデート] Amazon Redshift が Amazon Kinesis Data Streams からのストリーミング取り込みの並行スケーリングをサポートしたので試してみました
クラウド事業統括本部の石川です。Amazon Redshift のパッチ P203 以降で、Amazon Kinesis Data Streams に接続したストリーミングマテリアライズドビューの更新が並行スケーリングの対象になりましたので、実際に試してみました。
ストリーミング取り込みの更新処理をメインクラスターからスケーリングクラスタにオフロードできるようになるため、リアルタイムのデータ取り込みと BI・ETL などのワークロードがリソースを奪い合う構成で効果が期待できます。
Amazon Redshift ストリーミング取り込みとは
Amazon Redshift ストリーミング取り込み(Streaming Ingestion)は、Amazon Kinesis Data Streams や Amazon Managed Streaming for Apache Kafka から、低レイテンシーかつ高速にデータを取り込む機能です。
**Amazon S3 などの一時的な中継領域を経由せず、ストリームのデータが直接マテリアライズドビューに書き込まれます。**これにより外部データへ素早くアクセスでき、データアクセス時間とストレージコストの削減につながります。プロビジョンドクラスターと Amazon Redshift Serverless のワークグループのどちらでも、少数の SQL コマンドで構成できます。
アップデート内容
従来の課題
ストリーミングマテリアライズドビューには、明示的に REFRESH MATERIALIZED VIEW を実行する手動更新と、AUTO REFRESH を指定した自動更新の 2 通りがあります。
公式ドキュメントには、自動更新について次のように記載されています。
Auto refresh queries for a materialized view or views are treated as any other user workload.
(マテリアライズドビューの自動更新クエリは、他のユーザーワークロードと同様に扱われます)
つまり、ストリームからデータが到着するたびに走る更新クエリは、メインの Amazon Redshift クラスター上でユーザークエリと同じリソースを消費していました。取り込みが継続的に発生するストリーミングのワークロードでは、これがダッシュボードや ETL といった他ジョブとの競合要因になります。
何が変わったか
パッチ P203 以降、Amazon Redshift は Amazon Kinesis Data Streams に接続されたストリーミングマテリアライズドビューの更新について、並行スケーリングをサポートします。
主な変更点は以下のとおりです。
- Amazon Kinesis Data Streams 接続のストリーミングマテリアライズドビューの更新が、並行スケーリングの対象になった
- 並行スケーリングを有効にしていれば、ストリーミングワークロードが自動的にスケールする
- 更新処理がオフロードされることで、メインの Amazon Redshift クラスターを、他の優先度の高いワークロードの実行に振り向けられる
なお、並行スケーリングはこれまでも COPY・INSERT・DELETE・UPDATE・CREATE TABLE AS(CTAS)・VACUUM といった書き込みステートメントと、マテリアライズドビューの手動更新をサポートしていました。今回のアップデートは、そこに Amazon Kinesis Data Streams 接続のストリーミングマテリアライズドビューの更新が加わったかたちになります。
データの流れは以下のとおりです。
対応リージョン
AWS What's New では「Amazon Redshift が利用可能なすべての AWS リージョンで、すぐに利用開始できます」と案内されています。
一方で、並行スケーリング自体には公式ドキュメント上で対応リージョンの一覧が定義されており、東京リージョン(ap-northeast-1)や大阪リージョン(ap-northeast-3)を含む主要リージョンが対象です。今回の機能は並行スケーリングの上に成り立つため、利用前にご自身のリージョンが並行スケーリングの対応リージョンに含まれているかを、あわせて確認しておくと安全です。
料金への影響
今回の発表に、機能そのものの追加料金についての記載はありません。並行スケーリングの通常の料金体系が適用されます。
- メインクラスターが稼働している 24 時間ごとに、1 時間分の並行スケーリングクラスターのクレジットが加算される
- 無料クレジットは、アクティブなクラスターごとに最大 30 時間まで蓄積できる
- 無料クレジットを超えて並行スケーリングクラスターを使用した分は、秒単位のオンデマンド料金が課金される
ストリーミング取り込みは常時稼働のユースケースが多く、更新が頻繁に並行スケーリングクラスターへ回ると無料クレジットを消費しきる可能性があります。後述の max_concurrency_scaling_clusters による上限設定と、コストのモニタリングをセットで検討することをおすすめします。
やってみた
前提条件
- 検証リージョン: ap-northeast-1(東京)
- Amazon Redshift プロビジョンドクラスター: rg.large × 2 ノード(multi-node)
- Amazon Kinesis Data Streams: 1 シャード
並行スケーリングは、読み取りクエリであれば DC2 を含むノードタイプで利用できますが、書き込み操作は RG ノードと RA3 ノードのみでサポートされます。ストリーミングマテリアライズドビューの更新は、後述のとおり実体がマテリアライズドビュー内部テーブルへの INSERT であり書き込み操作にあたるため、この条件が適用されます。また並行スケーリングの対象となるにはシングルノードクラスターでないことが条件のため、2 ノード構成としました。
システム構成図

WLM で並行スケーリングを有効化する
Auto WLM を使い、concurrency_scaling を auto に設定します。あわせて並行スケーリングクラスターの最大数を 2 に指定します。
% aws redshift modify-cluster-parameter-group \
--region ap-northeast-1 \
--parameter-group-name blog-20260825-e9fba1eb-wlm \
--parameters '[
{"ParameterName":"wlm_json_configuration",
"ParameterValue":"[{\"auto_wlm\":true,\"concurrency_scaling\":\"auto\"},{\"short_query_queue\":true}]"},
{"ParameterName":"max_concurrency_scaling_clusters","ParameterValue":"2"}
]'
{
"ParameterGroupName": "blog-20260825-e9fba1eb-wlm",
"ParameterGroupStatus": "Your parameter group has been updated. If you changed only dynamic parameters, associated clusters are being modified now. If you changed static parameters, all updates, including dynamic parameters, will be applied when you reboot the associated clusters."
}
設定が反映されたことを確認します。
$ aws redshift describe-cluster-parameters \
--region ap-northeast-1 \
--parameter-group-name blog-20260825-e9fba1eb-wlm \
--query 'Parameters[?ParameterName==`wlm_json_configuration`||ParameterName==`max_concurrency_scaling_clusters`].{N:ParameterName,V:ParameterValue}' \
--output json
[
{
"N": "max_concurrency_scaling_clusters",
"V": "2"
},
{
"N": "wlm_json_configuration",
"V": "[{\"auto_wlm\":true,\"concurrency_scaling\":\"auto\"},{\"short_query_queue\":true}]"
}
]
クラスターを作成する
作成したIAM ロールを指定して、Amazon Redshift マルチノードクラスター(rg.large 2ノード)を作成します。
クラスターを作成する
$ aws redshift create-cluster \
--region ap-northeast-1 \
--cluster-identifier cs-kds-e9fba1eb \
--node-type rg.large \
--number-of-nodes 2 \
--cluster-type multi-node \
--master-username awsuser \
--master-user-password '<マスターパスワード>' \
--db-name dev \
--cluster-subnet-group-name blog-20260825-e9fba1eb-subnets \
--cluster-parameter-group-name blog-20260825-e9fba1eb-wlm \
--iam-roles arn:aws:iam::123456789012:role/blog-20260825-e9fba1eb-rs-streaming-role \
--no-publicly-accessible \
--tags Key=run_id,Value=blog-20260825-e9fba1eb \
Key=owner,Value=blog \
Key=purpose,Value=blog-verification
{
"Cluster": {
"ClusterIdentifier": "cs-kds-e9fba1eb",
"NodeType": "rg.large",
"ClusterStatus": "creating",
"ClusterAvailabilityStatus": "Modifying",
"MasterUsername": "awsuser",
"DBName": "dev",
"AutomatedSnapshotRetentionPeriod": 1,
"ManualSnapshotRetentionPeriod": -1,
"ClusterSecurityGroups": [],
"VpcSecurityGroups": [
{
"VpcSecurityGroupId": "sg-0548ccd1499215c3e",
"Status": "active"
}
],
"ClusterParameterGroups": [
{
"ParameterGroupName": "blog-20260825-e9fba1eb-wlm",
"ParameterApplyStatus": "in-sync"
}
],
"ClusterSubnetGroupName": "blog-20260825-e9fba1eb-subnets",
"VpcId": "vpc-0b0642bd5008af386",
"PreferredMaintenanceWindow": "thu:17:00-thu:17:30",
"PendingModifiedValues": {
"MasterUserPassword": "****"
},
"ClusterVersion": "1.0",
"AllowVersionUpgrade": true,
"NumberOfNodes": 2,
"PubliclyAccessible": false,
"Encrypted": true,
"Tags": [
{
"Key": "owner",
"Value": "blog"
},
{
"Key": "run_id",
"Value": "blog-20260825-e9fba1eb"
},
{
"Key": "purpose",
"Value": "blog-verification"
}
]
}
}
SQL の実行方法
クラスターを --no-publicly-accessible で作成したため、手元の端末から直接 TCP 接続できません。そこで Amazon Redshift Data API を使いました。Data API は AWS 管理のエンドポイント経由で動くため VPC 外からでも到達でき、--db-user による一時認証を使うのでマスターパスワードを保持する必要もありません。
Data API は非同期です。SQL を投げると即座にステートメント ID が返り、ステータスをポーリングしてから結果を取得します。
Data API によるSQL 実行
runsql() {
local sql="$1"
local id
# 1. SQL を投げてステートメント ID を受け取る
id=$(aws redshift-data execute-statement --region ap-northeast-1 \
--cluster-identifier cs-kds-e9fba1eb --database dev --db-user awsuser \
--sql "$sql" --query 'Id' --output text)
# 2. FINISHED / FAILED / ABORTED になるまで 2 秒間隔でポーリング
local st="" n=0
while [ $n -lt 150 ]; do
st=$(aws redshift-data describe-statement --region ap-northeast-1 --id "$id" --query 'Status' --output text)
case "$st" in FINISHED|FAILED|ABORTED) break;; esac
sleep 2; n=$((n+1))
done
# 3. 結果セットがあれば取得
aws redshift-data get-statement-result --region ap-northeast-1 --id "$id" --output json
}
get-statement-result は ColumnMetadata と Records に分かれた JSON を返すため、そのままでは読みづらく、以降に掲載する結果は列幅を揃えたテキストテーブルに整形しています。
パッチバージョンを確認する
今回の機能はパッチ P203 以降で利用できます。まずクラスターのバージョンを確認します。
SELECT version();
version
--------------------------------------------------------------------------------------------------------------------------
PostgreSQL 8.0.2 on i686-pc-linux-gnu, compiled by GCC gcc (GCC) 3.4.2 20041017 (Red Hat 3.4.2-6.fc3), Redshift 1.0.400034
(1 rows)
1.0.400034 は、公式ドキュメントのクラスターバージョン一覧でパッチ 203 に含まれるバージョンです。対象の機能が有効な環境であることが確認できました。
検証用データを投入する
センサーデータを模した JSON を 2,000 件投入しました。
2,000 件投入
import boto3, json, random
k = boto3.client('kinesis', region_name='ap-northeast-1')
random.seed(42)
devices = [f"device-{i:03d}" for i in range(1, 21)]
total = 0
for batch in range(4):
recs = []
for i in range(500):
n = batch * 500 + i
payload = {
"device_id": random.choice(devices),
"seq": n,
"temperature": round(random.uniform(15.0, 35.0), 2),
"humidity": round(random.uniform(30.0, 80.0), 2),
"status": random.choice(["ok", "ok", "ok", "warn", "error"]),
}
recs.append({"Data": json.dumps(payload).encode(), "PartitionKey": payload["device_id"]})
r = k.put_records(StreamName='blog-20260825-e9fba1eb-stream', Records=recs)
total += len(recs) - r.get('FailedRecordCount', 0)
print(f"batch {batch}: sent={len(recs)} failed={r.get('FailedRecordCount', 0)}")
print(f"TOTAL_OK={total}")
batch 0: sent=500 failed=0
batch 1: sent=500 failed=0
batch 2: sent=500 failed=0
batch 3: sent=500 failed=0
TOTAL_OK=2000
外部スキーマとストリーミングマテリアライズドビューを作成する
Amazon Kinesis Data Streams を参照する外部スキーマを作成します。
CREATE EXTERNAL SCHEMA kds
FROM KINESIS
IAM_ROLE 'arn:aws:iam::123456789012:role/blog-20260825-e9fba1eb-rs-streaming-role';
続いて、この外部スキーマ上にストリーミングマテリアライズドビューを作成します。AUTO REFRESH YES を指定して自動更新を有効にしました。
CREATE MATERIALIZED VIEW mv_sensor AUTO REFRESH YES AS
SELECT approximate_arrival_timestamp,
partition_key,
shard_id,
sequence_number,
JSON_PARSE(kinesis_data) AS payload
FROM kds."blog-20260825-e9fba1eb-stream";
公式ドキュメントでは、JSON_EXTRACT_PATH_TEXT を列ごとに繰り返し使うと抽出する列数だけ JSON が再パースされ取り込みレイテンシーが悪化するため、いったん JSON_PARSE で SUPER 型のまま取り込んでおき、あとから PartiQL で個々の値を取り出す方法が推奨されています。ここでもその方針に従いました。
取り込みを確認する
手動で更新をかけます。
REFRESH MATERIALIZED VIEW mv_sensor;
取り込み結果を確認します。
SELECT COUNT(*) AS ingested_rows FROM mv_sensor;
ingested_rows
-------------
2000
(1 rows)
SUPER 型に格納したペイロードを PartiQL で展開してみます。
SELECT partition_key,
payload.device_id::VARCHAR AS device_id,
payload.temperature::DECIMAL(5,2) AS temperature,
payload.status::VARCHAR AS status
FROM mv_sensor
ORDER BY sequence_number
LIMIT 5;
partition_key | device_id | temperature | status
--------------+------------+-------------+-------
device-004 | device-004 | 15.50 | ok
device-005 | device-005 | 29.73 | error
device-003 | device-003 | 26.81 | ok
device-007 | device-007 | 19.65 | error
device-007 | device-007 | 29.32 | warn
(5 rows)
問題なく取り込めています。
WLM の設定を確認する
並行スケーリングが有効になっているか確認します。
SELECT service_class, num_query_tasks, name, concurrency_scaling
FROM stv_wlm_service_class_config
WHERE service_class >= 5
ORDER BY service_class;
service_class | num_query_tasks | name | concurrency_scaling
--------------+-----------------+------------------------------------------------------------------+---------------------
5 | 1 | Service class for super user | off
14 | 6 | Short query queue | off
15 | 0 | Service class for vacuum/analyze | off
100 | -1 | Default queue | auto
(4 rows)
Default queue(service_class 100)が auto になっています。Auto WLM のため num_query_tasks は -1(動的)です。
更新がどこで実行されたかを確認する
ここまでのクエリがメインクラスターと並行スケーリングクラスターのどちらで実行されたかを確認します。
SELECT query,
concurrency_scaling_status,
TRIM(SUBSTRING(querytxt,1,45)) AS query_text,
starttime
FROM stl_query
WHERE userid > 1
ORDER BY starttime DESC
LIMIT 12;
query | concurrency_scaling_status | query_text | starttime
------+----------------------------+-----------------------------------------------+---------------------------
3186 | 6 | SELECT service_class, num_query_tasks, name, | 2026-08-25 05:34:45.958598
3173 | 0 | SELECT partition_key, payload.device_id::VARC | 2026-08-25 05:34:23.805455
3168 | 0 | SELECT COUNT(*) AS ingested_rows FROM mv_sens | 2026-08-25 05:34:21.386898
3163 | 0 | INSERT INTO "public"."mv_tbl__mv_sensor__0" | 2026-08-25 05:34:17.099753
3157 | 0 | INSERT INTO "public"."mv_tbl__mv_sensor__0" | 2026-08-25 05:34:10.342018
3153 | 0 | INSERT INTO "public"."mv_tbl__mv_sensor__0" | 2026-08-25 05:34:08.428581
3148 | 5 | CREATE MATERIALIZED VIEW mv_sensor AUTO REFRE | 2026-08-25 05:34:00.92326
(7 rows)
ストリーミングマテリアライズドビューの更新は、裏側にある実体テーブルへの INSERT INTO "public"."mv_tbl__mv_sensor__0" として記録されることがわかります。この時点では concurrency_scaling_status はすべて 0、つまりメインクラスターでの実行です。
公式ドキュメントでは concurrency_scaling_status = 1 が並行スケーリングクラスターでの実行を示すと説明されています。上記に現れている 5 や 6 といった値については、ドキュメント上に一覧の記載を見つけられませんでした。ドキュメントのサンプルクエリも CASE WHEN concurrency_scaling_status = 1 THEN ... ELSE 'main cluster' END という形で、1 かそれ以外かで判定しています。
更新の履歴は SVL_MV_REFRESH_STATUS でも確認できます。
SELECT TRIM(db_name) AS db, TRIM(mv_name) AS mv, status,
TRIM(refresh_type) AS refresh_type, starttime
FROM svl_mv_refresh_status
ORDER BY starttime DESC LIMIT 10;
db | mv | status | refresh_type | starttime
----+-----------+----------------------------------------------------------------------------+--------------+---------------------------
dev | mv_sensor | Refresh successfully updated MV incrementally. Stream returned no new data | Manual | 2026-08-25 05:34:16.976869
dev | mv_sensor | Refresh successfully updated MV incrementally. Stream returned no new data | Auto | 2026-08-25 05:34:10.16889
dev | mv_sensor | Refresh successfully updated MV incrementally | Auto | 2026-08-25 05:34:08.246986
(3 rows)
興味深いことに、こちらを見ると先に Auto(自動更新)が動いてデータを取り込んでおり、後から実行した手動 REFRESH のときには「Stream returned no new data」となっています。AUTO REFRESH YES を指定していたため、明示的に REFRESH するまでもなく取り込みが完了していたわけです。
AUTO REFRESH の動作を確認する
自動更新の挙動をもう少し確かめます。手動 REFRESH を実行せずに、ストリームへ 3,000 件を追加投入しました。
ストリームへ 3,000 件を追加投入
import boto3, json, random
k = boto3.client('kinesis', region_name='ap-northeast-1')
random.seed(7)
devices = [f"device-{i:03d}" for i in range(1, 21)]
total = 0
for b in range(6):
recs = []
for i in range(500):
n = 100000 + b*500 + i
p = {"device_id": random.choice(devices), "seq": n,
"temperature": round(random.uniform(15, 35), 2),
"humidity": round(random.uniform(30, 80), 2),
"status": random.choice(["ok", "warn", "error"])}
recs.append({"Data": json.dumps(p).encode(), "PartitionKey": p["device_id"]})
r = k.put_records(StreamName='blog-20260825-e9fba1eb-stream', Records=recs)
total += 500 - r.get('FailedRecordCount', 0)
print(f"ADDED_OK={total}")
ADDED_OK=3000
しばらく待ってから件数を確認します。
SELECT COUNT(*) AS rows_now FROM mv_sensor;
rows_now
--------
5000
(1 rows)
AUTO REFRESH YES の指定どおり、追加分が自動的に取り込まれています。
負荷をかけて並行スケーリングの発動を試みる
並行スケーリングはキューが滞留したときに発動します。そこで負荷生成用のテーブルを用意しました。ストリーミングマテリアライズドビュー同士をクロスジョインして 100 万行のテーブルを作ります。
なお、Amazon Redshift の CTAS は IF NOT EXISTS をサポートしていないため、以下のように記述します。
CREATE TABLE load_seed AS
SELECT a.sequence_number AS sn,
b.partition_key AS pk,
a.payload.temperature::DECIMAL(5,2) AS temp
FROM mv_sensor a CROSS JOIN mv_sensor b
LIMIT 1000000;
SELECT COUNT(*) AS seed_rows FROM load_seed;
seed_rows
---------
1000000
(1 rows)
このテーブルに対する 13 秒程度かかるクエリを 20 本、同時に投入します。Data API は非同期で即座に返るため、シェルの & で並列投入できます。
13 秒程度かかるクエリを 20 本、同時に投入
HEAVY='SELECT a.pk, b.pk AS pk2, COUNT(*) AS c FROM load_seed a JOIN load_seed b ON a.sn = b.sn GROUP BY 1,2 ORDER BY 3 DESC LIMIT 20;'
for i in $(seq 1 20); do
aws redshift-data execute-statement --region ap-northeast-1 \
--cluster-identifier cs-kds-e9fba1eb --database dev --db-user awsuser \
--sql "${HEAVY}" --statement-name "cs-probe-${i}" --query 'Id' --output text >/dev/null &
done
wait
投入直後から STV_INFLIGHT をポーリングしましたが、並行スケーリングクラスターでの実行は観測できませんでした。
poll 1 [05:41:35Z]: cs_cluster=0 inflight_total=14
poll 2 [05:41:43Z]: cs_cluster=0 inflight_total=13
poll 3 [05:41:51Z]: cs_cluster=0 inflight_total=13
poll 4 [05:41:59Z]: cs_cluster=0 inflight_total=13
poll 5 [05:42:07Z]: cs_cluster=0 inflight_total=13
poll 6 [05:42:15Z]: cs_cluster=0 inflight_total=13
実行中クエリが 13〜14 本で頭打ちになっており、残りはキューで待っているはずです。処理が落ち着いてから集計しても、並行スケーリングクラスターでの実行は 0 件でした。
SELECT CASE WHEN q.concurrency_scaling_status = 1
THEN 'concurrency scaling cluster' ELSE 'main cluster' END AS run_on,
COUNT(*) AS queries,
SUM(ROUND(w.total_queue_time::NUMERIC/1000000,2)) AS queue_secs,
SUM(ROUND(w.total_exec_time::NUMERIC/1000000,2)) AS exec_secs
FROM stl_query q JOIN stl_wlm_query w USING (userid, query)
WHERE q.userid > 1 AND q.starttime > DATEADD(minute, -10, GETDATE())
GROUP BY 1 ORDER BY 1;
run_on | queries | queue_secs | exec_secs
-------------+---------+------------+----------
main cluster | 55 | 1.31 | 64.39
(1 rows)
キュー待ち時間の合計が 1.31 秒しかありません。Auto WLM が同時実行数を動的に調整し、負荷を捌き切っていたようです。
手動 WLM でスロット数を絞る
Auto WLM ではキューが滞留しづらいため、手動 WLM に切り替えてスロット数を明示的に 2 まで絞りました。
% aws redshift modify-cluster-parameter-group \
--region ap-northeast-1 \
--parameter-group-name blog-20260825-e9fba1eb-wlm \
--parameters '[
{"ParameterName":"wlm_json_configuration",
"ParameterValue":"[{\"query_group\":[],\"query_group_wild_card\":0,\"user_group\":[],\"user_group_wild_card\":0,\"concurrency_scaling\":\"auto\",\"query_concurrency\":2,\"max_execution_time\":0,\"memory_percent_to_use\":100},{\"short_query_queue\":false}]"},
{"ParameterName":"max_concurrency_scaling_clusters","ParameterValue":"2"}
]'
auto_wlm の切り替えは静的パラメータのため、再起動が必要です。
% aws redshift reboot-cluster \
--region ap-northeast-1 \
--cluster-identifier cs-kds-e9fba1eb \
--query 'Cluster.{Status:ClusterStatus,PG:ClusterParameterGroups[0].ParameterApplyStatus}' \
--output json
{
"Status": "available",
"PG": "pending-reboot"
}
しばらく待つと in-sync になります。
{
"Status": "available",
"PG": "in-sync"
}
スロット数を確認します。
SELECT service_class, num_query_tasks AS slots, TRIM(name) AS name, concurrency_scaling
FROM stv_wlm_service_class_config WHERE service_class >= 5 ORDER BY service_class;
service_class | slots | name | concurrency_scaling
--------------+-------+----------------------------------+---------------------
5 | 1 | Service class for super user | off
6 | 2 | Default queue | auto
15 | 0 | Service class for vacuum/analyze | off
(3 rows)
Default queue のスロットが 2 になりました。Auto WLM のときは service_class 100 でしたが、手動 WLM では 6 になっています。
手動 WLM で並行スケーリングの発動を観測する(リトライ)
ストリームへ 3,000 件を追加投入して自動更新を発火させたうえで、同じクエリを 10 本投入しました。スロットが 2 本なので 8 本がキュー待ちになる想定です。
ストリームへ 3,000 件を追加投入
for i in $(seq 1 10); do
aws redshift-data execute-statement --region ap-northeast-1 \
--cluster-identifier cs-kds-e9fba1eb --database dev --db-user awsuser \
--sql "${HEAVY}" --statement-name "mwlm-${i}" --query 'Id' --output text >/dev/null &
done
wait
ポーリングした結果です。
poll 1 [06:01:15Z]: 確認クエリ自体が status=STARTED (キュー待ち)
poll 2 [06:01:24Z]: cs_cluster=0 inflight=2
poll 3 [06:01:34Z]: 確認クエリ自体が status=STARTED (キュー待ち)
poll 4 [06:01:43Z]: 確認クエリ自体が status=STARTED (キュー待ち)
poll 5 [06:01:53Z]: cs_cluster=5 inflight=7
poll 6 [06:02:03Z]: cs_cluster=0 inflight=2
poll 7 [06:02:13Z]: cs_cluster=0 inflight=1
poll 8 [06:02:22Z]: cs_cluster=0 inflight=1
poll 5 で、実行中の 7 本のうち 5 本が並行スケーリングクラスターに回っていることが観測できました。なお poll 1、3、4 では確認用のクエリ自体がキューで待たされ、結果が返ってきていません。監視クエリも同じ WLM キューを通るため、メインクラスターが詰まると監視そのものが刺さります。
実行場所を集計する
処理が落ち着いてから改めて集計します。
SELECT CASE WHEN q.concurrency_scaling_status = 1
THEN 'concurrency scaling cluster' ELSE 'main cluster' END AS run_on,
COUNT(*) AS queries,
SUM(ROUND(w.total_queue_time::NUMERIC/1000000,2)) AS queue_secs,
SUM(ROUND(w.total_exec_time::NUMERIC/1000000,2)) AS exec_secs
FROM stl_query q JOIN stl_wlm_query w USING (userid, query)
WHERE q.userid > 1 AND q.starttime > DATEADD(minute, -6, GETDATE())
GROUP BY 1 ORDER BY 1;
run_on | queries | queue_secs | exec_secs
----------------------------+---------+------------+----------
concurrency scaling cluster | 5 | 226.95 | 51.27
main cluster | 17 | 209.53 | 122.15
(2 rows)
並行スケーリングの利用実績は SVCS_CONCURRENCY_SCALING_USAGE で確認できます。
SELECT start_time, end_time, queries, usage_in_seconds
FROM svcs_concurrency_scaling_usage ORDER BY start_time DESC LIMIT 5;
start_time | end_time | queries | usage_in_seconds
---------------------------+----------------------------+---------+-----------------
2026-08-25 06:01:41.6675 | 2026-08-25 06:01:52.372099 | 5 | 11
2026-08-25 05:42:55.814568 | 2026-08-25 05:42:57.20169 | 1 | 2
2026-08-25 05:42:28.936983 | 2026-08-25 05:42:38.691135 | 4 | 10
2026-08-25 05:42:13.815229 | 2026-08-25 05:42:25.166128 | 5 | 12
(4 rows)
ここで気づいたことがあります。05:42 台に 3 回の利用実績が記録されています。これは Auto WLM のまま負荷をかけていた時間帯です。つまり、発動していなかったのではなく、発動していたのを観測できていなかったということでした。並行スケーリングクラスターの稼働は 10 秒前後と短く、STV_INFLIGHT のポーリングでは捉えきれていませんでした。
先ほど「並行スケーリングクラスターでの実行は 0 件」と集計した際に対象としていた時間帯(DATEADD(minute, -10, GETDATE()))は、この発動より前だったことになります。
ストリーミングマテリアライズドビューの更新はどこで実行されたか
本題です。ストリーミングマテリアライズドビューの取り込みクエリが、どちらのクラスターで実行されたかを確認します。
SELECT query,
concurrency_scaling_status AS cs_status,
TRIM(SUBSTRING(querytxt,1,42)) AS query_text,
starttime
FROM stl_query
WHERE userid > 1 AND querytxt ILIKE '%mv_tbl__mv_sensor%'
ORDER BY starttime DESC LIMIT 20;
query | cs_status | query_text | starttime
-------+-----------+-------------------------------------------+---------------------------
409377 | 0 | INSERT INTO "public"."mv_tbl__mv_sensor__ | 2026-08-25 06:03:11.812773
409365 | 0 | INSERT INTO "public"."mv_tbl__mv_sensor__ | 2026-08-25 06:02:40.959074
409341 | 0 | INSERT INTO "public"."mv_tbl__mv_sensor__ | 2026-08-25 06:02:10.896532
409321 | 0 | INSERT INTO "public"."mv_tbl__mv_sensor__ | 2026-08-25 06:01:51.65466
409279 | 0 | INSERT INTO "public"."mv_tbl__mv_sensor__ | 2026-08-25 06:01:08.968543
3803 | 0 | INSERT INTO "public"."mv_tbl__mv_sensor__ | 2026-08-25 05:44:28.778121
3797 | 0 | INSERT INTO "public"."mv_tbl__mv_sensor__ | 2026-08-25 05:43:58.888783
3784 | 0 | INSERT INTO "public"."mv_tbl__mv_sensor__ | 2026-08-25 05:43:26.83914
3765 | 1 | INSERT INTO "public"."mv_tbl__mv_sensor__ | 2026-08-25 05:42:55.78917
3722 | 1 | INSERT INTO "public"."mv_tbl__mv_sensor__ | 2026-08-25 05:41:53.838842
3609 | 0 | INSERT INTO "public"."mv_tbl__mv_sensor__ | 2026-08-25 05:41:22.78361
3604 | 0 | INSERT INTO "public"."mv_tbl__mv_sensor__ | 2026-08-25 05:40:51.913177
3577 | 0 | INSERT INTO "public"."mv_tbl__mv_sensor__ | 2026-08-25 05:40:23.294678
3573 | 0 | INSERT INTO "public"."mv_tbl__mv_sensor__ | 2026-08-25 05:40:20.843386
3458 | 0 | INSERT INTO "public"."mv_tbl__mv_sensor__ | 2026-08-25 05:39:18.923916
3444 | 0 | INSERT INTO "public"."mv_tbl__mv_sensor__ | 2026-08-25 05:38:47.910329
3347 | 0 | INSERT INTO "public"."mv_tbl__mv_sensor__ | 2026-08-25 05:38:16.806437
3331 | 0 | INSERT INTO "public"."mv_tbl__mv_sensor__ | 2026-08-25 05:37:48.332626
3327 | 0 | INSERT INTO "public"."mv_tbl__mv_sensor__ | 2026-08-25 05:37:45.898323
3163 | 0 | INSERT INTO "public"."mv_tbl__mv_sensor__ | 2026-08-25 05:34:17.099753
(20 rows)
query 3722(05:41:53)と query 3765(05:42:55)の 2 件で cs_status = 1 になっています。ストリーミングマテリアライズドビューの自動更新が、並行スケーリングクラスターで実行されたことが確認できました。今回のアップデートの動作そのものです。
時刻を見ると、SVCS_CONCURRENCY_SCALING_USAGE に記録されていた 05:42 台の発動と一致しています。混雑していない 05:40 台前半や 06:01 以降の更新は、いずれもメインクラスター(cs_status = 0)で実行されています。負荷に応じて振り分けられている様子がうかがえます。
最後に、更新回数と最終的な行数を確認します。
SELECT TRIM(refresh_type) AS refresh_type, COUNT(*) AS n, MAX(starttime) AS latest
FROM svl_mv_refresh_status GROUP BY 1 ORDER BY 2 DESC;
refresh_type | n | latest
-------------+----+---------------------------
Auto | 21 | 2026-08-25 06:03:11.681697
Manual | 1 | 2026-08-25 05:34:16.976869
(2 rows)
SELECT COUNT(*) AS total_rows FROM mv_sensor;
total_rows
----------
12000
(1 rows)
自動更新 21 回、手動更新 1 回で、合計 12,000 行が取り込まれました。
考察
今回の検証で確認できた点と、注意しておきたい点を整理します。
確認できたこと
- パッチ P203 環境において、Amazon Kinesis Data Streams 接続のストリーミングマテリアライズドビューの自動更新が、並行スケーリングクラスターで実行される
- 常に並行スケーリングへ回るわけではなく、メインクラスターが混雑している時間帯に限って振り分けられる
- 自動更新の実体は
STL_QUERYに実体テーブルへのINSERTとして記録されるため、concurrency_scaling_statusで実行場所を追跡できる
観測にあたっての注意点
並行スケーリングクラスターの稼働は 10 秒前後と短く、STV_INFLIGHT を数秒おきにポーリングする方法では取りこぼします。実際、今回も一度は「発動しなかった」と誤って判断しかけました。事後に SVCS_CONCURRENCY_SCALING_USAGE と STL_QUERY を突き合わせるほうが確実です。
また、確認用のクエリ自体も WLM のキューに入るため、メインクラスターが詰まっている状況では結果が返ってきません。監視は別経路、たとえば Amazon CloudWatch の ConcurrencyScalingSeconds や ConcurrencyScalingActiveClusters メトリクスと併用するのが現実的です。
コスト面
並行スケーリングは、メインクラスターが稼働している 24 時間ごとに 1 時間分のクレジットが加算され、アクティブなクラスターごとに最大 30 時間まで蓄積されます。無料クレジットを超えた分は秒単位のオンデマンド料金が課金されます。
ストリーミング取り込みは常時稼働のユースケースが多く、更新が頻繁に並行スケーリングへ回る構成では無料クレジットを消費しきる可能性があります。max_concurrency_scaling_clusters による上限設定とコストのモニタリングをセットで検討しておくとよさそうです。
制限事項
- 書き込み操作の並行スケーリングは RG ノードと RA3 ノードのみでサポートされます。これは今回の Amazon Kinesis Data Streams 連携に固有の制限ではなく、COPY・INSERT・DELETE・UPDATE・CTAS・VACUUM などを含む書き込み操作全般に共通する条件です(読み取りクエリであれば DC2 でも並行スケーリングを利用できます)
- シングルノードクラスターは並行スケーリングの対象外です
最後に
Amazon Kinesis Data Streams に接続したストリーミングマテリアライズドビューの更新が、パッチ P203 以降で並行スケーリングの対象になりました。実際に rg.large × 2 ノードのクラスターで検証したところ、メインクラスターが混雑している時間帯に、自動更新が並行スケーリングクラスターへオフロードされることを確認できました。
ニアリアルタイム分析の基盤では、取り込み処理とダッシュボードのクエリがリソースを奪い合う構成になりがちです。並行スケーリングを有効にしておくことで、取り込みの遅延とクエリの応答時間を両立させやすくなります。
すでにストリーミング取り込みを運用されている方は、クラスターのパッチバージョンとノードタイプを確認したうえで、WLM キューの並行スケーリングを有効化し、SVCS_CONCURRENCY_SCALING_USAGE で実際の使用状況と無料クレジットの消費ペースを見てみてはいかがでしょうか。










