30歳からのプログラミング

30歳無職から独学でプログラミングを開始した人間の記録。

Datastream で Amazon Aurora から BigQuery にストリーミング転送する

Datastream は Google Cloud が提供しているサービスのひとつで、データのストリーミング転送を実現できる。
バッチ処理などに比べてデータの同期をリアルタイム性高く行うことができる。

この記事では、細かい設定方法などについては扱わず、 Datastream の考え方や仕組み、ストリーミング転送を有効にした際の挙動などについて述べていく。

Datastream は様々な転送元、転送先を選べるが、今回は Amazon Aurora の PostgreSQL から BigQuery への転送を行う。

Aurora インスタンスと Datastream の疎通

まずは Aurora と Datastream の疎通について述べる。
それぞれ別のクラウドサービスが提供するサービスなので、両者を疎通させる経路をインターネット上に作る必要がある。
そのための方法を Datastream は複数用意しているが、この記事では、 AWS に踏み台サーバを用意しそれを経由して疎通する方法を採用することにする。以降、それを前提とした説明となる。

まず、踏み台サーバとして EC2 インスタンスを用意する。そして踏み台サーバから Aurora インスタンスにアクセスできるようにしておく。そして、インターネットゲートウェイが設置された VPC ネットワークに踏み台サーバを配置することで、外部から踏み台サーバに、そしてそれを経由して Aurora インスタンスにアクセスできるようにする。
Datastream は SSH で踏み台サーバにアクセスする。そのため、セキュリティグループの設定で、 Datastream が使う IP アドレスが TCP のポート22(SSH で使うポート)に接続できるようにしておく必要がある。Datastream が使う IP アドレスは以下のページに掲載されている。
https://docs.cloud.google.com/datastream/docs/ip-allowlists-and-regions

これで、 Aurora インスタンスと Datastream が情報をやり取りできるようになる。

Aurora 側の設定

PostgreSQL には「論理レプリケーション」という仕組みがあり、これを使うことでストリーミング転送できるようになる。
論理レプリケーションとは、「このテーブルにこの内容で INSERT した」「このテーブルのこのレコードを DELETE した」という変更内容を記録し、記録した内容を他のソフトウェアに配信する仕組みのこと。
変更内容を受け取るソフトウェアを「コンシューマ」と呼ぶ。この記事の文脈でいうと、 Datastream がコンシューマとなる。

Amazon Aurora の PostgreSQL の場合、rds.logical_replicationというパラメータを1にすると論理レプリケーションが有効になる。

しかしこれだけでは不十分で、論理レプリケーションを動かすためのオブジェクトを用意する必要がある。
まず、 Publication 。これは、どのテーブルの変更を配信するのかを定義するオブジェクト。
そして、 Replication Slot 。これは、コンシューマの情報を管理するためのオブジェクト。このオブジェクトに、コンシューマがどこまで変更内容を受け取ったのかを記録する。そうすることで、まだ処理されていない変更内容が失われないようにする。
この 2 つを用意しなければならない。

そして、コンシューマ(今回のケースだと Datastream )が使用する、 PostgreSQL のユーザーも用意する必要がある。そのユーザーは当然、必要な権限を持っている必要がある。

Datastream 側の設定

「接続プロファイル」と「ストリーム」という、 2 種類のリソースを作成する。

接続プロファイルはその名の通り、対象とどのように接続するのかを管理するリソースであり、転送先と転送元それぞれに用意する。
今回のケースでは転送先は BigQuery だが、その場合、転送先を表現する接続プロファイルには、これといって書くことはない。Datastream も BigQuery も Google Cloud のサービスであるため、 Google Cloud が内部でよしなにしてくれる。利用者が何かを意識する必要がない。
一方、転送元である Aurora インスタンスとの接続を表現する接続プロファイルには、様々な情報を書く必要がある。対象のインスタンスのエンドポイントやポート番号、踏み台サーバの IP アドレスやポート番号、使用する PostgreSQL ユーザー、など。踏み台サーバに SSH 接続するための秘密鍵の情報も、接続プロファイルに登録する。

最後に、ストリームというリソースを作成し、それを有効にする。
ストリームには、 2 つの接続プロファイルの他、使用する Publication や Replication Slot の名前など、ストリーミング転送そのものの設定を記述する。また、転送対象のテーブルも指定する。ストリームで指定されており、かつ Publication に含まれているテーブルが、転送対象となる。
ストリームに設定するパラメータのうち、特に重要なもののひとつがdata_freshness。この値が、転送先の各テーブルのmax_stalenessの初期値になる。後述するが、max_stalenessの値によって、転送先にデータが反映される頻度が変わる。

ストリームが適切に作成され有効になると、 Datastream が論理レプリケーションのコンシューマとなり、論理レプリケーションが動作し始める。
Datastream を利用する、というのは、ストリームを作成しそれを有効にする、ということであり、ここまでの準備はストリームを作るための準備であった、と整理できるかもしれない。

ストリームを有効にした際の挙動

ストリームの設定で「バックフィル」というものを有効にすると、ストリームを有効にした時点で転送元にあった全てのレコードが転送先のテーブルに作られる。
そしてそのあとは、転送元のテーブルに対して INSERT, UPDATE, DELETE が行われる度に、その内容が転送先のテーブルにも反映される。

ではどの程度のタイミング、頻度で反映されるのか。それに影響を与えるのが、先程軽く触れたmax_stalenessである。
これは BigQuery の各テーブルが持つパラメータで、例えばこの値が「15分」である場合、少なくとも 15 分前までの変更内容がクエリ結果に反映されることが保証される。逆に言えば、転送元での変更がクエリ結果に現れるまでに、最大で 15 分程度の遅れが生じ得る、ということである。反映のタイミングを厳密に指定したり、任意のタイミングで手動で反映させたり、といったことはできない。
この値を小さくするほどデータが新鮮になるが、その分 BigQuery 側で変更を反映するための処理が頻繁に実行されるようになり、金銭的コストが増える。

レコードに対する操作ではなく、ALTER TABLEによるカラムの追加や削除を行うとどうなるのか。

追加にしろ削除にしろ、ALTER TABLEしただけでは、その変更は BigQuery には反映されない。

例えばordersテーブルにnoteカラムを追加したとする。

id status amount note
1 active 1200 NULL
2 pending 800 NULL
3 cancelled 0 NULL

しかし BigQuery 側にはnoteカラムは存在しないままである。

id status amount
1 active 1200
2 pending 800
3 cancelled 0

このあとに INSERT, UPDATE, DELETE などを行ったときに初めて、 BigQuery 側にも反映される。
例えば転送元のテーブルで、idが1のレコードと2のレコードのnoteに値を入れる UPDATE を実行したとする。
それによって転送元のテーブルが以下の状態になった場合、 BigQuery 側のテーブルも同じ内容になる。

id status amount note
1 active 1200 Initial registration
2 pending 800 Awaiting payment
3 cancelled 0 NULL

カラム削除については、 BigQuery に反映されることはない。
例えば、転送元でnoteカラムを削除したとする。

id status amount
1 active 1200
2 pending 800
3 cancelled 0

しかし BigQuery 側は、変わらず以下の状態である。

id status amount note
1 active 1200 Initial registration
2 pending 800 Awaiting payment
3 cancelled 0 NULL

では転送元でレコードの編集を行うとどうなるか。
idが2のレコードを UPDATE して、以下のようにしたとする(amountを800から999にしている)。

id status amount
1 active 1200
2 pending 999
3 cancelled 0

そうすると BigQuery 側は以下のようになる。

id status amount note
1 active 1200 Initial registration
2 pending 999 NULL
3 cancelled 0 NULL

idが2のレコードが更新され、そのレコードのnoteはNULLになる。
noteというカラムは残り続け、変更があったレコード以外のレコードについては値は変わらない。上記の例だとidが1のレコードのnoteカラムの値は、このレコードに対する操作が発生しない限りInitial registrationのまま変化しない。

Kube Resource Orchestrator を使って Kubernetes のリソースを抽象化して提供する

Kube Resource Orchestrator ( kro ) は Kubernetes のリソースを管理するためのツールで、複数のリソースをひとつのまとまりとして定義したり、リソースとリソースの依存関係を表現したりすることができる。
そして kro は定義された内容に基づいて、適切な順番でリソースを作成してくれる。
kro を上手く使うことで、 Kubernetes クラスタの複雑さを隠蔽し、利用者にとって使いやすいインターフェースを提供することができる。

kro は Kubernetes 組み込みの機能ではないため、使用するためには自分でセットアップする必要がある。だが Amazon EKS はマネージドサービスとして kro を提供しており、それを使えば比較的簡単にセットアップできる。

この記事では、 kro の概要や使い方、 Amazon EKS で kro を使用する方法、活用例、などについて述べていく。

この記事の内容は Kubernetes の1.34で動作確認した。

EKS Capability for kro のセットアップ

まず、 Amazon EKS のクラスタで kro を使えるようにするためのセットアップを行う。
kro そのものの説明は次のセクションから行うので、セットアップに興味がなければ読み飛ばしても問題ない。

Amazon EKS は EKS Capability という機能を提供している。これを使うことで、便利な各種ソフトウェアをクラスタにインストールすることなく利用できる。
対象のソフトウェアに kro も含まれているため( EKS Capability for kro )、それを有効にすることで kro を使えるようになる。Amazon EKS が提供しているマネージドサービスであるため、 kro 自体の管理をしなくて済む。

対象のクラスタで EKS Capability for kro が有効になっているかは、以下のコマンドで確認できる。
$ aws eks list-capabilities --cluster-name <CLUSTER_NAME> --region <REGION>

今回はap-northeast-1リージョンに作ったmy-clusterという名前のクラスタを対象にする。

$ aws eks list-capabilities --cluster-name my-cluster --region ap-northeast-1
{
    "capabilities": []
}

my-clusterでは EKS Capability はまだ何も有効になっていないので、有効にするための作業を進めていく。

$ aws eks create-capabilityというコマンドを実行すれば有効にできるが、このコマンドには AWS の IAM ロールを渡す必要があるので、まずは IAM ロールを作る必要がある。
以下の手順で進める。

  1. IAM ロールに付与する信頼ポリシーを記述した JSON ファイルを用意する
  2. KROCapabilityRoleという名前の IAM ロールを作る
  3. $ aws eks create-capabilityを実行する

以下の内容が書かれたkro-trust-policy.jsonを用意する。

{
  "Version": "2012-10-17",
  "Statement": [
    {
      "Effect": "Allow",
      "Principal": {
        "Service": "capabilities.eks.amazonaws.com"
      },
      "Action": [
        "sts:AssumeRole",
        "sts:TagSession"
      ]
    }
  ]
}

このように書くことで、この信頼ポリシーをアタッチして作られた IAM ロールをcapabilities.eks.amazonaws.com、つまり EKS Capability が使えるようになる。

以下のコマンドを実行して、KROCapabilityRoleという名前の IAM ロールを作る。
$ aws iam create-role --role-name KROCapabilityRole --assume-role-policy-document file://kro-trust-policy.json

そして$ aws eks create-capabilityを実行する。
$ aws eks create-capability --region ap-northeast-1 --cluster-name my-cluster --capability-name my-kro --type KRO --role-arn arn:aws:iam::<AWS_ACCOUNT_ID>:role/KROCapabilityRole --delete-propagation-policy RETAIN

これで EKS Capability for kro が有効になる。

$ aws eks list-capabilities --cluster-name my-cluster --region ap-northeast-1
{
    "capabilities": [
        {
            "capabilityName": "my-kro",
            "arn": "arn:aws:eks:ap-northeast-1:<AWS_ACCOUNT_ID>:capability/my-cluster/kro/my-kro/88cf84b6-e30c-9b47-4816-b09be9448c44",
            "type": "KRO",
            "status": "ACTIVE",
            "version": "0.9.2-eks-2",
            "createdAt": "2026-06-27T23:13:37.748000+09:00",
            "modifiedAt": "2026-06-27T23:13:37.748000+09:00"
        }
    ]
}

kro の基本的な使い方

RGD を作る

Kubernetes は拡張性が高く、 Pod や Job といった組み込みのリソースを利用するだけでなく、リソースを新しく定義してそれを利用することができる。
kro は ResourceGraphDefinition ( RGD ) というリソースを提供しており、それを使って「リソースを作るリソース」を作ることができる。

具体例を見たほうが分かりやすいので示す。以下は、JobTriggerという RGD を定義している例。

apiVersion: kro.run/v1alpha1
kind: ResourceGraphDefinition
metadata:
  name: job-trigger
spec:
  schema:
    apiVersion: v1alpha1
    kind: JobTrigger
    spec:
      name: string
      value: string
  resources:
    - id: job
      template:
        apiVersion: batch/v1
        kind: Job
        metadata:
          name: ${schema.spec.name}
        spec:
          template:
            spec:
              containers:
                - name: echo
                  image: busybox
                  command: ["echo"]
                  args: ["${schema.spec.value}"]
              restartPolicy: Never

まず記法を説明したあと、この内容を apply すると何が起こるのかを見ていく。

spec.schemaで、この RGD がどのようなものであるか、を定義している。この例では、JobTriggerはnameとvalueという 2 つのフィールドを持つ、ということを定義している。
spec.resourcesは、JobTriggerから作られるリソースについて定義している。

resourcesは配列であり複数のリソースを定義できるが、上記の例ではひとつの Job のみを定義している。
リソースの書き方は何も特別なことはなく、そのリソースに応じた書き方をすればよい。kro を使っているからといって、必要なフィールドが変わったりはしない。
${schema.***}という記法でschemaのフィールドの値を参照することができる。例えばJobTriggerではargs: ["${schema.spec.value}"]と書いているが、これにより、valueがfooのJobTriggerを作成したら、そこから作成される Job のargsは["foo"]になる。つまりこの Job は「JobTriggerのvalueに渡された値を標準出力する Job」ということになる。

この内容を apply すると、JobTriggerがクラスタに登録される。

$ kubectl get resourcegraphdefinition
NAME          APIVERSION   KIND         STATE    READY   AGE
job-trigger   v1alpha1     JobTrigger   Active   True    3s

これでJobTriggerのインスタンスを作成できるようになる。

JobTrigger のインスタンスを作成する

マニフェストファイルに以下の内容を書き apply する。

apiVersion: kro.run/v1alpha1
kind: JobTrigger
metadata:
  name: test
spec:
  name: my-job
  value: hello

そうすると、testという名前のJobTriggerが作られる。
そしてこのJobTriggerから、my-jobという名前の Job が作られ、その Job はecho helloを実行するはずである。

確認してみると、 JobTrigger は確かに作られているが、STATEがERRORになっている。そして Job は作られていない。

$ kubectl get jobtrigger
NAME   STATE   READY   AGE
test   ERROR   False   73s

$ kubectl get job
No resources found in default namespace.

testの詳細を調べてみると、 kro が使っている IAM ロールKROCapabilityRoleが job に関する権限を持っておらず、それでエラーになっていることがわかる。

$ kubectl get jobtrigger test -o json | jq -r '[.status.conditions[] | select(.status == "False") | .message] | unique | .[]'
resource reconciliation failed: jobs.batch "my-job" is forbidden: User "arn:aws:sts::<AWS_ACCOUNT_ID>:assumed-role/KROCapabilityRole/KRO" cannot get resource "jobs" in API group "batch" in the namespace "default"

KROCapabilityRoleは$ aws eks create-capabilityの引数で渡した IAM ロール。EKS Capability for kro において kro は、この IAM ロールを使ってクラスタにアクセスし操作を行う。そのためこの IAM ロールでは認可されていない操作を行おうとすると、このようにエラーになってしまう。
そのため次は、KROCapabilityRoleに権限を付与する。

まずkro-job-managerという ClusterRole を作る。Job に関する一通りの権限を持たせる。

apiVersion: rbac.authorization.k8s.io/v1
kind: ClusterRole
metadata:
  name: kro-job-manager
rules:
  - apiGroups: ["batch"]
    resources: ["jobs"]
    verbs: ["create", "get", "list", "watch", "delete", "patch", "update"]

そしてkro-job-managerをKROCapabilityRoleにバインドする。

apiVersion: rbac.authorization.k8s.io/v1
kind: ClusterRoleBinding
metadata:
  name: kro-job-manager-binding
subjects:
  - kind: User
    name: "arn:aws:sts::<AWS_ACCOUNT_ID>:assumed-role/KROCapabilityRole/KRO"
    apiGroup: rbac.authorization.k8s.io
roleRef:
  kind: ClusterRole
  name: kro-job-manager
  apiGroup: rbac.authorization.k8s.io

そうすることで、 kro がこのクラスタで Job のcreateなどを行えるようになる。

apply すると、testのSTATEがACTIVEになっており、 Job も作られている。

$ kubectl get jobtrigger
NAME   STATE    READY   AGE
test   ACTIVE   True    12m

$ kubectl get job
NAME     STATUS     COMPLETIONS   DURATION   AGE
my-job   Complete   1/1           6s         42s

そして Job のログを見るとhelloが出力されている。

$ kubectl logs job/my-job
hello

kro の活用例

kro を上手く使うことで、 Kubernetes クラスタの複雑さを隠蔽し、適切に抽象化されたプラットフォームを提供することができる。

例えば、先程のJobTriggerを作成する権限をクラスタの利用者に配り、 Kubernetes 組み込みの Job を作る権限は渡さないようにする。そうすると、クラスタ管理者がJobTriggerによって定義した Job は実行できるが、それ以外の任意の Job を自由に作成・実行することはできない、という状況を作ることができる。これにより、利用者の利便性と安全性を両立させることができる。

クラスタ利用者が複雑な設定を理解し記述しなければならない、という状況の解決策にもなり得る。
例えば、ウェブアプリケーションをクラスタ上に構築するためには多くの Kubernetes リソースを記述しなければならない状況だとする。それらのリソースはどのようなものなのか、リソース間の依存関係はどうなっているのか、ということを記述した RGD を用意することで、ウェブアプリケーションを構築したい人はただその RGD を使えばよくなる。RGD で定義している個々のリソースについて知らなくて済むようになる。

複雑さを隠蔽しウェブアプリケーションを構築しやすくする RGD 、のサンプルが以下。

apiVersion: kro.run/v1alpha1
kind: ResourceGraphDefinition
metadata:
  name: simple-app
spec:
  schema:
    apiVersion: v1alpha1
    kind: SimpleApp
    spec:
      appName: string
      image: string
      port: integer | default=80
    status:
      availableReplicas: ${myDeployment.status.availableReplicas}
  resources:
    - id: myDeployment
      template:
        apiVersion: apps/v1
        kind: Deployment
        metadata:
          name: ${schema.spec.appName}
        spec:
          replicas: 1
          selector:
            matchLabels:
              app: ${schema.spec.appName}
          template:
            metadata:
              labels:
                app: ${schema.spec.appName}
            spec:
              containers:
                - name: app
                  image: ${schema.spec.image}
                  ports:
                    - containerPort: ${schema.spec.port}
    - id: myService
      template:
        apiVersion: v1
        kind: Service
        metadata:
          name: ${schema.spec.appName}-svc
        spec:
          type: LoadBalancer
          selector: ${myDeployment.spec.selector.matchLabels}
          ports:
            - port: 80
              targetPort: ${schema.spec.port}
              protocol: TCP

この RGD SimpleAppは、 Deployment と Service を定義している。
JobTriggerにはなかった記法として、${myDeployment.***}が使われている。これはidがmyDeploymentであるリソースのspecやstatusを参照している。この記法を使うことで、リソース間の依存関係を表現できる。
そして kro はこの依存関係を理解し、それに基づいて適切な順番でリソースを作成してくれる。

SimpleAppが用意されていることで、クラスタ利用者は Deployment や Service について意識することなく、ウェブアプリケーションをクラスタ上に構築できる。
ただ以下のようにSimpleAppのインスタンスを作成すればよい。

apiVersion: kro.run/v1alpha1
kind: SimpleApp
metadata:
  name: my-simple-app
spec:
  appName: my-app
  image: nginx:latest

SimpleAppを実際に使うためにはKROCapabilityRoleに以下の権限を追加しておく必要があるので、これも忘れずに追記しておく。

   - apiGroups: ["batch"]
     resources: ["jobs"]
     verbs: ["create", "get", "list", "watch", "delete", "patch", "update"]
+  - apiGroups: ["apps"]
+    resources: ["deployments"]
+    verbs: ["create", "get", "list", "watch", "delete", "patch", "update"]
+  - apiGroups: [""]
+    resources: ["services"]
+    verbs: ["create", "get", "list", "watch", "delete", "patch", "update"]

そのうえで apply すると、my-simple-appが作られる。

$ kubectl get simpleapp
NAME            STATE    READY   AGE
my-simple-app   ACTIVE   True    2s

そしてSimpleAppの定義に基づいて、my-appという Deployment とmy-app-svcという Service が作られる。

$ kubectl get deploy my-app
NAME     READY   UP-TO-DATE   AVAILABLE   AGE
my-app   1/1     1            1           34s

$ kubectl get service my-app-svc
NAME         TYPE           CLUSTER-IP       EXTERNAL-IP                                                                   PORT(S)        AGE
my-app-svc   LoadBalancer   172.20.166.209   a67319105608041cc934c65d679ea56a-365706414.ap-northeast-1.elb.amazonaws.com   80:31506/TCP   45s

これでウェブアプリケーションが構築され、公開される。

$ curl -I http://a67319105608041cc934c65d679ea56a-365706414.ap-northeast-1.elb.amazonaws.com
HTTP/1.1 200 OK
Server: nginx/1.31.2
Date: Sat, 27 Jun 2026 16:16:24 GMT
Content-Type: text/html
Content-Length: 896
Last-Modified: Wed, 17 Jun 2026 14:40:35 GMT
Connection: keep-alive
ETag: "6a32b1e3-380"
Accept-Ranges: bytes

SimpleAppの例では Deployment と Service が 1 つずつあるだけなのでありがたみは薄いかもしれないが、リソースが多くなり依存関係が複雑になればなるほど、 kro による抽象化の恩恵は増していく。