Snowflake Openflow Connector for PostgreSQL を SPCS で動かし、RDS for PostgreSQLに接続してみた
データ事業本部の笠原です。
先日SnowflakeのOpenflow Connector for PostgreSQLを使って、RDS for PostgreSQLのデータをCDC (Change Data Capture) でSnowflakeに転送できるか、試しました。
この時は、OpenflowのデプロイモデルのうちBYOCを採用し、自身のAWSアカウント上のEKS上でOpenflowを動かしてみました。
今回はもう一つのデプロイモデルであるSPCSを採用して、Snowflake上でOpenflowを動かして、RDS for PostgreSQLに接続できるか試してみました。
今回の構成
今回の構成は以下の通りです。

OpenflowのランタイムはSPCS (Snowpark Container Services) 上で動作しています。SPCSはSnowflake内で完全に自己完結するサービスであるため、デプロイと管理が容易になります。
RDS for PostgreSQLへの接続ですが、今回はPrivateLink接続ではなく、SnowflakeのStandard/Enterpriseライセンスでも使えるように、インターネット経由での接続を想定して構成しました。
RDSをパブリックネットワークからアクセス可能にする方法も考えられますが、今回はNLBを経由することでSPCSから接続できるようにしてみました。
NLBのセキュリティグループにて、インバウンドルールにはSnowflakeのegress CIDRを設定して、IPアドレス制限を行います。また、NLBのターゲットグループにはRDSエンドポイントからIPアドレスを名前解決して設定するLambda関数を用意します。
構築手順
1. ネットワーク環境
今回はパブリックサブネット2つ、プライベートサブネット2つのシンプルな構成にしています。また、以降で使用するセキュリティグループも事前に準備してます。
RDS for PostgreSQLインスタンスに対してSQLを実行する踏み台EC2は、今回このテンプレートで作成しています。
なお、NLBはこのテンプレートではなく、別のテンプレートで作成します。
ネットワーク環境構築Cfnテンプレート:01-network.yaml
AWSTemplateFormatVersion: '2010-09-09'
Parameters:
ProjectName:
Type: String
Default: openflow-pg-nlb
VpcCidr:
Type: String
Default: 10.1.0.0/16
PublicSubnet1Cidr:
Type: String
Default: 10.1.0.0/24
PublicSubnet2Cidr:
Type: String
Default: 10.1.1.0/24
PrivateSubnet1Cidr:
Type: String
Default: 10.1.10.0/24
PrivateSubnet2Cidr:
Type: String
Default: 10.1.11.0/24
BastionInstanceType:
Type: String
Default: t3.micro
LatestAmiId:
Type: AWS::SSM::Parameter::Value<AWS::EC2::Image::Id>
Default: /aws/service/ami-amazon-linux-latest/al2023-ami-kernel-default-x86_64
Resources:
Vpc:
Type: AWS::EC2::VPC
Properties:
CidrBlock: !Ref VpcCidr
EnableDnsSupport: true
EnableDnsHostnames: true
Tags:
- Key: Name
Value: !Sub '${ProjectName}-vpc'
InternetGateway:
Type: AWS::EC2::InternetGateway
Properties:
Tags:
- Key: Name
Value: !Sub '${ProjectName}-igw'
IgwAttachment:
Type: AWS::EC2::VPCGatewayAttachment
Properties:
VpcId: !Ref Vpc
InternetGatewayId: !Ref InternetGateway
PublicSubnet1:
Type: AWS::EC2::Subnet
Properties:
VpcId: !Ref Vpc
CidrBlock: !Ref PublicSubnet1Cidr
AvailabilityZone: !Select [0, !GetAZs '']
MapPublicIpOnLaunch: true
Tags:
- Key: Name
Value: !Sub '${ProjectName}-public-1'
PublicSubnet2:
Type: AWS::EC2::Subnet
Properties:
VpcId: !Ref Vpc
CidrBlock: !Ref PublicSubnet2Cidr
AvailabilityZone: !Select [1, !GetAZs '']
MapPublicIpOnLaunch: true
Tags:
- Key: Name
Value: !Sub '${ProjectName}-public-2'
PrivateSubnet1:
Type: AWS::EC2::Subnet
Properties:
VpcId: !Ref Vpc
CidrBlock: !Ref PrivateSubnet1Cidr
AvailabilityZone: !Select [0, !GetAZs '']
Tags:
- Key: Name
Value: !Sub '${ProjectName}-private-1'
PrivateSubnet2:
Type: AWS::EC2::Subnet
Properties:
VpcId: !Ref Vpc
CidrBlock: !Ref PrivateSubnet2Cidr
AvailabilityZone: !Select [1, !GetAZs '']
Tags:
- Key: Name
Value: !Sub '${ProjectName}-private-2'
PublicRouteTable:
Type: AWS::EC2::RouteTable
Properties:
VpcId: !Ref Vpc
Tags:
- Key: Name
Value: !Sub '${ProjectName}-rt-public'
DefaultRoute:
Type: AWS::EC2::Route
DependsOn: IgwAttachment
Properties:
RouteTableId: !Ref PublicRouteTable
DestinationCidrBlock: 0.0.0.0/0
GatewayId: !Ref InternetGateway
PublicSubnet1RtAssoc:
Type: AWS::EC2::SubnetRouteTableAssociation
Properties:
SubnetId: !Ref PublicSubnet1
RouteTableId: !Ref PublicRouteTable
PublicSubnet2RtAssoc:
Type: AWS::EC2::SubnetRouteTableAssociation
Properties:
SubnetId: !Ref PublicSubnet2
RouteTableId: !Ref PublicRouteTable
PrivateRouteTable:
Type: AWS::EC2::RouteTable
Properties:
VpcId: !Ref Vpc
Tags:
- Key: Name
Value: !Sub '${ProjectName}-rt-private'
PrivateSubnet1RtAssoc:
Type: AWS::EC2::SubnetRouteTableAssociation
Properties:
SubnetId: !Ref PrivateSubnet1
RouteTableId: !Ref PrivateRouteTable
PrivateSubnet2RtAssoc:
Type: AWS::EC2::SubnetRouteTableAssociation
Properties:
SubnetId: !Ref PrivateSubnet2
RouteTableId: !Ref PrivateRouteTable
NlbSecurityGroup:
Type: AWS::EC2::SecurityGroup
Properties:
GroupDescription: internet-facing NLB - inbound 5432 from Snowflake egress IPs (added later).
VpcId: !Ref Vpc
SecurityGroupEgress:
- IpProtocol: -1
CidrIp: 0.0.0.0/0
Description: Allow all outbound (forward to RDS target + health checks).
Tags:
- Key: Name
Value: !Sub '${ProjectName}-nlb-sg'
RdsSecurityGroup:
Type: AWS::EC2::SecurityGroup
Properties:
GroupDescription: RDS PostgreSQL (private) - inbound 5432 from NLB SG and bastion only.
VpcId: !Ref Vpc
SecurityGroupIngress:
- IpProtocol: tcp
FromPort: 5432
ToPort: 5432
SourceSecurityGroupId: !Ref NlbSecurityGroup
Description: PostgreSQL from the internet-facing NLB nodes.
- IpProtocol: tcp
FromPort: 5432
ToPort: 5432
SourceSecurityGroupId: !Ref BastionSecurityGroup
Description: PostgreSQL from bastion (psql admin / postgres setup SQL).
Tags:
- Key: Name
Value: !Sub '${ProjectName}-rds-sg'
BastionSecurityGroup:
Type: AWS::EC2::SecurityGroup
Properties:
GroupDescription: Bastion host - outbound only (SSM Session Manager, dnf, psql to RDS).
VpcId: !Ref Vpc
SecurityGroupEgress:
- IpProtocol: -1
CidrIp: 0.0.0.0/0
Description: Allow all outbound.
Tags:
- Key: Name
Value: !Sub '${ProjectName}-bastion-sg'
LambdaSecurityGroup:
Type: AWS::EC2::SecurityGroup
Properties:
GroupDescription: IP-sync Lambda - outbound only (ELB + Logs via interface endpoints).
VpcId: !Ref Vpc
SecurityGroupEgress:
- IpProtocol: -1
CidrIp: 0.0.0.0/0
Description: Allow all outbound (HTTPS to interface endpoints, DNS).
Tags:
- Key: Name
Value: !Sub '${ProjectName}-lambda-sg'
VpcEndpointSecurityGroup:
Type: AWS::EC2::SecurityGroup
Properties:
GroupDescription: VPC interface endpoints - inbound 443 from the Lambda SG.
VpcId: !Ref Vpc
SecurityGroupIngress:
- IpProtocol: tcp
FromPort: 443
ToPort: 443
SourceSecurityGroupId: !Ref LambdaSecurityGroup
Description: HTTPS from the IP-sync Lambda.
Tags:
- Key: Name
Value: !Sub '${ProjectName}-vpce-sg'
ElbInterfaceEndpoint:
Type: AWS::EC2::VPCEndpoint
Properties:
VpcId: !Ref Vpc
ServiceName: !Sub 'com.amazonaws.${AWS::Region}.elasticloadbalancing'
VpcEndpointType: Interface
PrivateDnsEnabled: true
SubnetIds:
- !Ref PrivateSubnet1
- !Ref PrivateSubnet2
SecurityGroupIds:
- !Ref VpcEndpointSecurityGroup
LogsInterfaceEndpoint:
Type: AWS::EC2::VPCEndpoint
Properties:
VpcId: !Ref Vpc
ServiceName: !Sub 'com.amazonaws.${AWS::Region}.logs'
VpcEndpointType: Interface
PrivateDnsEnabled: true
SubnetIds:
- !Ref PrivateSubnet1
- !Ref PrivateSubnet2
SecurityGroupIds:
- !Ref VpcEndpointSecurityGroup
BastionRole:
Type: AWS::IAM::Role
Properties:
AssumeRolePolicyDocument:
Version: '2012-10-17'
Statement:
- Effect: Allow
Principal:
Service: ec2.amazonaws.com
Action: sts:AssumeRole
ManagedPolicyArns:
- arn:aws:iam::aws:policy/AmazonSSMManagedInstanceCore
Tags:
- Key: Name
Value: !Sub '${ProjectName}-bastion-role'
BastionInstanceProfile:
Type: AWS::IAM::InstanceProfile
Properties:
Roles:
- !Ref BastionRole
BastionInstance:
Type: AWS::EC2::Instance
Properties:
InstanceType: !Ref BastionInstanceType
ImageId: !Ref LatestAmiId
IamInstanceProfile: !Ref BastionInstanceProfile
SubnetId: !Ref PublicSubnet1
SecurityGroupIds:
- !Ref BastionSecurityGroup
UserData:
Fn::Base64: !Sub |
#!/bin/bash
# PostgreSQL client (psql) for running postgres/*.sql against RDS.
dnf install -y postgresql16 || dnf install -y postgresql15
Tags:
- Key: Name
Value: !Sub '${ProjectName}-bastion'
Outputs:
VpcId:
Value: !Ref Vpc
Export:
Name: !Sub '${ProjectName}-VpcId'
PublicSubnet1Id:
Value: !Ref PublicSubnet1
Export:
Name: !Sub '${ProjectName}-PublicSubnet1Id'
PublicSubnet2Id:
Value: !Ref PublicSubnet2
Export:
Name: !Sub '${ProjectName}-PublicSubnet2Id'
PrivateSubnet1Id:
Value: !Ref PrivateSubnet1
Export:
Name: !Sub '${ProjectName}-PrivateSubnet1Id'
PrivateSubnet2Id:
Value: !Ref PrivateSubnet2
Export:
Name: !Sub '${ProjectName}-PrivateSubnet2Id'
NlbSecurityGroupId:
Description: Attach Snowflake egress IPs here (aws/scripts/update-snowflake-egress-sg.sh --sg-id <this>).
Value: !Ref NlbSecurityGroup
Export:
Name: !Sub '${ProjectName}-NlbSecurityGroupId'
RdsSecurityGroupId:
Value: !Ref RdsSecurityGroup
Export:
Name: !Sub '${ProjectName}-RdsSecurityGroupId'
BastionSecurityGroupId:
Value: !Ref BastionSecurityGroup
Export:
Name: !Sub '${ProjectName}-BastionSecurityGroupId'
LambdaSecurityGroupId:
Value: !Ref LambdaSecurityGroup
Export:
Name: !Sub '${ProjectName}-LambdaSecurityGroupId'
BastionInstanceId:
Description: Connect with `aws ssm start-session --target <id>`.
Value: !Ref BastionInstance
Export:
Name: !Sub '${ProjectName}-BastionInstanceId'
2. データベース (RDS for PostgreSQL) 環境
データソースとなるDBを構築します。前回と変わりません。
- rds.logical_replication = 1
- これにより wal_level=logical になります。
新規作成インスタンスには作成時から適用されますが、念のため SHOW wal_level; で確認しましょう。
- これにより wal_level=logical になります。
- PubliclyAccessible = false
- NLB経由でアクセスするため、パブリックアクセスはしない設定にします。
DB環境構築Cfnテンプレート:02-rds-postgres.yaml
AWSTemplateFormatVersion: '2010-09-09'
Parameters:
ProjectName:
Type: String
Default: openflow-pg-nlb
EngineVersion:
Type: String
Default: '16.14'
ParameterGroupFamily:
Type: String
Default: postgres16
DBInstanceClass:
Type: String
Default: db.t4g.micro
AllocatedStorage:
Type: Number
Default: 20
DBName:
Type: String
Default: appdb
MasterUsername:
Type: String
Default: postgres
Resources:
DBSubnetGroup:
Type: AWS::RDS::DBSubnetGroup
Properties:
DBSubnetGroupDescription: !Sub '${ProjectName} RDS subnet group (private)'
SubnetIds:
- Fn::ImportValue: !Sub '${ProjectName}-PrivateSubnet1Id'
- Fn::ImportValue: !Sub '${ProjectName}-PrivateSubnet2Id'
Tags:
- Key: Name
Value: !Sub '${ProjectName}-db-subnet-group'
DBParameterGroup:
Type: AWS::RDS::DBParameterGroup
Properties:
Description: !Sub '${ProjectName} - logical replication enabled'
Family: !Ref ParameterGroupFamily
Parameters:
rds.logical_replication: '1'
Tags:
- Key: Name
Value: !Sub '${ProjectName}-pg-params'
DBInstance:
Type: AWS::RDS::DBInstance
DeletionPolicy: Delete
UpdateReplacePolicy: Delete
Properties:
DBInstanceIdentifier: !Sub '${ProjectName}-pg'
Engine: postgres
EngineVersion: !Ref EngineVersion
DBInstanceClass: !Ref DBInstanceClass
AllocatedStorage: !Ref AllocatedStorage
StorageType: gp3
DBName: !Ref DBName
MasterUsername: !Ref MasterUsername
ManageMasterUserPassword: true # stores the master password in Secrets Manager
DBSubnetGroupName: !Ref DBSubnetGroup
DBParameterGroupName: !Ref DBParameterGroup
VPCSecurityGroups:
- Fn::ImportValue: !Sub '${ProjectName}-RdsSecurityGroupId'
PubliclyAccessible: false # private; reachable only via the NLB + bastion
MultiAZ: false
BackupRetentionPeriod: 1
DeletionProtection: false
Tags:
- Key: Name
Value: !Sub '${ProjectName}-pg'
Outputs:
DBEndpointAddress:
Description: >-
RDS endpoint host. Use as the RDS_ENDPOINT env for the IP-sync Lambda (03 stack).
NOTE: clients connect via the NLB DNS name, not this host. Its TLS cert CN is this
endpoint, so connect with sslmode=require (verify-full against the NLB name fails).
Value: !GetAtt DBInstance.Endpoint.Address
Export:
Name: !Sub '${ProjectName}-DBEndpointAddress'
DBEndpointPort:
Value: !GetAtt DBInstance.Endpoint.Port
Export:
Name: !Sub '${ProjectName}-DBEndpointPort'
DBName:
Value: !Ref DBName
MasterUserSecretArn:
Description: Secrets Manager ARN holding the master username/password.
Value: !GetAtt DBInstance.MasterUserSecret.SecretArn
Export:
Name: !Sub '${ProjectName}-MasterUserSecretArn'
3. NLB環境
今回利用するNLBと、NLBのターゲットグループにRDSエンドポイントのIPを設定するLambda関数を構築します。
NLB環境構築Cfnテンプレート 03-nlb-ipsync.yaml
AWSTemplateFormatVersion: '2010-09-09'
Parameters:
ProjectName:
Type: String
Default: openflow-pg-nlb
SyncScheduleExpression:
Type: String
Default: rate(1 minute)
LogRetentionDays:
Type: Number
Default: 14
Resources:
Nlb:
Type: AWS::ElasticLoadBalancingV2::LoadBalancer
Properties:
Name: !Sub '${ProjectName}-nlb'
Type: network
Scheme: internet-facing
IpAddressType: ipv4
SecurityGroups:
- Fn::ImportValue: !Sub '${ProjectName}-NlbSecurityGroupId'
Subnets:
- Fn::ImportValue: !Sub '${ProjectName}-PublicSubnet1Id'
- Fn::ImportValue: !Sub '${ProjectName}-PublicSubnet2Id'
LoadBalancerAttributes:
- Key: load_balancing.cross_zone.enabled
Value: 'true'
- Key: deletion_protection.enabled
Value: 'false'
Tags:
- Key: Name
Value: !Sub '${ProjectName}-nlb'
TargetGroup:
Type: AWS::ElasticLoadBalancingV2::TargetGroup
Properties:
Name: !Sub '${ProjectName}-tg'
TargetType: ip # RDS のプライベートIPは、IP-sync用のLambdaにて設定します。
Protocol: TCP
Port: 5432
VpcId:
Fn::ImportValue: !Sub '${ProjectName}-VpcId'
HealthCheckProtocol: TCP
HealthCheckPort: traffic-port
HealthCheckIntervalSeconds: 30
HealthyThresholdCount: 3
UnhealthyThresholdCount: 3
TargetGroupAttributes:
- Key: preserve_client_ip.enabled
Value: 'false'
- Key: deregistration_delay.timeout_seconds
Value: '30'
Tags:
- Key: Name
Value: !Sub '${ProjectName}-tg'
Listener:
Type: AWS::ElasticLoadBalancingV2::Listener
Properties:
LoadBalancerArn: !Ref Nlb
Protocol: TCP
Port: 5432
DefaultActions:
- Type: forward
TargetGroupArn: !Ref TargetGroup
IpSyncLogGroup:
Type: AWS::Logs::LogGroup
Properties:
LogGroupName: !Sub '/aws/lambda/${ProjectName}-nlb-ipsync'
RetentionInDays: !Ref LogRetentionDays
IpSyncRole:
Type: AWS::IAM::Role
Properties:
AssumeRolePolicyDocument:
Version: '2012-10-17'
Statement:
- Effect: Allow
Principal:
Service: lambda.amazonaws.com
Action: sts:AssumeRole
ManagedPolicyArns:
- arn:aws:iam::aws:policy/service-role/AWSLambdaVPCAccessExecutionRole
Policies:
- PolicyName: nlb-target-sync
PolicyDocument:
Version: '2012-10-17'
Statement:
- Effect: Allow
Action:
- elasticloadbalancing:RegisterTargets
- elasticloadbalancing:DeregisterTargets
Resource: !Ref TargetGroup
- Effect: Allow
Action:
- elasticloadbalancing:DescribeTargetHealth
- elasticloadbalancing:DescribeTargetGroups
Resource: '*'
Tags:
- Key: Name
Value: !Sub '${ProjectName}-ipsync-role'
IpSyncFunction:
Type: AWS::Lambda::Function
DependsOn: IpSyncLogGroup
Properties:
FunctionName: !Sub '${ProjectName}-nlb-ipsync'
Runtime: python3.12
Handler: index.handler
Role: !GetAtt IpSyncRole.Arn
Timeout: 60
MemorySize: 128
ReservedConcurrentExecutions: 1
VpcConfig:
SubnetIds:
- Fn::ImportValue: !Sub '${ProjectName}-PrivateSubnet1Id'
- Fn::ImportValue: !Sub '${ProjectName}-PrivateSubnet2Id'
SecurityGroupIds:
- Fn::ImportValue: !Sub '${ProjectName}-LambdaSecurityGroupId'
Environment:
Variables:
RDS_ENDPOINT:
Fn::ImportValue: !Sub '${ProjectName}-DBEndpointAddress'
TARGET_GROUP_ARN: !Ref TargetGroup
PORT: '5432'
Code:
ZipFile: |
import os
import socket
import boto3
TG_ARN = os.environ["TARGET_GROUP_ARN"]
HOST = os.environ["RDS_ENDPOINT"]
PORT = int(os.environ.get("PORT", "5432"))
elbv2 = boto3.client("elbv2")
def resolve_ipv4(host):
infos = socket.getaddrinfo(host, PORT, family=socket.AF_INET,
type=socket.SOCK_STREAM)
return sorted({info[4][0] for info in infos})
def current_targets():
resp = elbv2.describe_target_health(TargetGroupArn=TG_ARN)
return sorted({d["Target"]["Id"]
for d in resp.get("TargetHealthDescriptions", [])})
def handler(event, context):
desired = resolve_ipv4(HOST)
existing = current_targets()
to_add = [ip for ip in desired if ip not in existing]
to_remove = [ip for ip in existing if ip not in desired]
if to_add:
try:
elbv2.register_targets(
TargetGroupArn=TG_ARN,
Targets=[{"Id": ip, "Port": PORT} for ip in to_add],
)
print("registered: %s" % to_add)
except Exception as e:
print("register failed for %s: %s" % (to_add, e))
if to_remove and desired:
try:
elbv2.deregister_targets(
TargetGroupArn=TG_ARN,
Targets=[{"Id": ip, "Port": PORT} for ip in to_remove],
)
print("deregistered: %s" % to_remove)
except Exception as e:
print("deregister failed for %s: %s" % (to_remove, e))
result = {"host": HOST, "desired": desired,
"added": to_add, "removed": to_remove}
print(result)
return result
Tags:
- Key: Name
Value: !Sub '${ProjectName}-nlb-ipsync'
SyncScheduleRule:
Type: AWS::Events::Rule
Properties:
Name: !Sub '${ProjectName}-nlb-ipsync-schedule'
ScheduleExpression: !Ref SyncScheduleExpression
State: ENABLED
Targets:
- Id: ipsync
Arn: !GetAtt IpSyncFunction.Arn
SyncSchedulePermission:
Type: AWS::Lambda::Permission
Properties:
FunctionName: !Ref IpSyncFunction
Action: lambda:InvokeFunction
Principal: events.amazonaws.com
SourceArn: !GetAtt SyncScheduleRule.Arn
Outputs:
NlbDnsName:
Value: !GetAtt Nlb.DNSName
Export:
Name: !Sub '${ProjectName}-NlbDnsName'
TargetGroupArn:
Value: !Ref TargetGroup
Export:
Name: !Sub '${ProjectName}-TargetGroupArn'
IpSyncFunctionName:
Value: !Ref IpSyncFunction
ちなみに、データ量が多く、初回のスナップショット取得に時間がかかると予想される場合、リスナーのアイドルタイムアウト時間を延長してください。
aws elbv2 modify-listener-attributes \
--listener-arn "$LISTENER_ARN" \
--attributes Key=tcp.idle_timeout.seconds,Value=6000
なお、NLBのターゲットグループの宛先の候補として、RDS Proxyを使う案も考えました。調査したところ、以下のRDS for PostgreSQLを使用したRDS Proxyの追加の制限事項により、今回RDS Proxyの利用は断念しました。
RDS Proxy は現在、ストリーミングレプリケーションモードをサポートしていません。
4. IP sync用Lambdaを実行し、ターゲットIPを登録
スケジュール実行でも登録されますが、即時反映と確認のため手動実行します。
## 実行
aws lambda invoke \
--function-name openflow-pg-nlb-nlb-ipsync /dev/stdout
## 確認
aws elbv2 describe-target-health \
--target-group-arn "$TG_ARN" \
--query "TargetHealthDescriptions[].{ip:Target.Id,state:TargetHealth.State}" \
--output table
5. PostgreSQL側設定
PostgreSQL側に、データソースとなるデータベース・スキーマ・テーブルをサンプル用に作成しておきます。
RDS for PostgreSQLへログインするマスターユーザパスワードは、SecretsManagerに格納されているので、確認ください。
前回同様SSMセッションマネージャーのポートフォワーディングを使って、ローカルPCからpsqlコマンドを実行しても良いですし、踏み台EC2上でpsqlコマンドを実行しても構いません。
SSMセッションマネージャーのポートフォワーディングを使う場合は、以下のようにssmのstart-sessionコマンドを実行します。
export BASTION_ID=<踏み台EC2のインスタンスID>
export RDS_ENDPOINT=<DBエンドポイント>
aws ssm start-session --region "ap-northeast-1" --target "$BASTION_ID" \
--document-name AWS-StartPortForwardingSessionToRemoteHost \
--parameters "{\"host\":[\"$RDS_ENDPOINT\"],\"portNumber\":[\"5432\"],\"localPortNumber\":[\"55432\"]}"
上記コマンドを実行しているターミナルは閉じずに、別のターミナルを立ち上げて、以下のコマンドを実行します。
psql "host=localhost port=55432 dbname=appdb user=pgadmin sslmode=require" \
-f postgres/01-test-data.sql
-f オプションで、SQLファイルを指定して実行します。
実行するSQLファイルの中身は以下の通りです。
このSQLファイルを実行することで、テストデータを投入しています。
テストデータ投入SQL 01-test-data.sql
CREATE TABLE IF NOT EXISTS public.customers (
id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
name TEXT NOT NULL,
email TEXT NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
INSERT INTO public.customers (name, email) VALUES
('Alice Tanaka', 'alice@example.com'),
('Bob Suzuki', 'bob@example.com'),
('Carol Yamada', 'carol@example.com');
CREATE TABLE IF NOT EXISTS public.orders (
id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
customer_id BIGINT NOT NULL REFERENCES public.customers(id),
amount NUMERIC(12,2) NOT NULL,
status TEXT NOT NULL DEFAULT 'NEW',
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
INSERT INTO public.orders (customer_id, amount, status) VALUES
(1, 1200.00, 'NEW'),
(1, 450.50, 'PAID'),
(2, 9800.00, 'NEW');
-- 件数check
SELECT 'customers' AS table, count(*) FROM public.customers
UNION ALL
SELECT 'orders' AS table, count(*) FROM public.orders;
次に、CDCを行うための論理レプリケーション設定の確認等を行います。
この際、Openflow Connector用に接続するPostgreSQLユーザのパスワードは適宜変更して登録してください。
SQLクエリ内の <CHANGE_ME_STRONG_PASSWORD> が該当箇所ですので、この値を変更してください。
論理レプリケーション設定確認、パブリケーション設定 02-logical-replication-setup.sql
-- Checl logical replication
SHOW wal_level;
SHOW max_replication_slots;
SHOW max_wal_senders;
-- PUBLICATION
CREATE PUBLICATION snowflake_pub WITH (publish_via_partition_root = true);
ALTER PUBLICATION snowflake_pub ADD TABLE public.customers, public.orders;
-- user for openflow connector
CREATE ROLE openflow_repl WITH LOGIN PASSWORD '<CHANGE_ME_STRONG_PASSWORD>';
GRANT rds_replication TO openflow_repl;
-- Grant for Snapshot / CDC
GRANT CONNECT ON DATABASE public TO openflow_repl;
GRANT USAGE ON SCHEMA public TO openflow_repl;
GRANT SELECT ON ALL TABLES IN SCHEMA public TO openflow_repl;
ALTER DEFAULT PRIVILEGES IN SCHEMA public GRANT SELECT ON TABLES TO openflow_repl;
-- Check
SELECT pubname, schemaname, tablename
FROM pg_publication_tables
WHERE pubname = 'snowflake_pub'
ORDER BY tablename;
6. Snowflake側設定
続いて、Snowflake側の設定に移ります。
Snowsight上のワークシートから、以下のSQLを実行します。
この際、クエリ内の <your_user> は実際にSnowflakeを作業するユーザ名、 <NLB_DNS> には作成したNLBのDNS名に、それぞれ置き換えてください。
Openflow有効化 03-openflow-deployment-setup.sql
USE ROLE ACCOUNTADMIN;
-- Create Openflow Admin Role
CREATE ROLE IF NOT EXISTS OPENFLOW_ADMIN;
-- Grant for Openflow
GRANT CREATE OPENFLOW DATA PLANE INTEGRATION ON ACCOUNT TO ROLE OPENFLOW_ADMIN;
GRANT CREATE OPENFLOW RUNTIME INTEGRATION ON ACCOUNT TO ROLE OPENFLOW_ADMIN;
GRANT CREATE COMPUTE POOL ON ACCOUNT TO ROLE OPENFLOW_ADMIN;
-- Openflow実行に使用するウェアハウス作成
CREATE WAREHOUSE IF NOT EXISTS OPENFLOW_WH
WAREHOUSE_SIZE = 'XSMALL'
AUTO_SUSPEND = 60
AUTO_RESUME = TRUE
INITIALLY_SUSPENDED = TRUE;
-- 自分の作業ユーザに付与
GRANT ROLE OPENFLOW_ADMIN TO USER <YOUR_USER>;
-- ランタイム/コネクタへの OAuth ログインは、作業ユーザの DEFAULT_ROLE を使う。
-- 作業ユーザのロールを `OPENFLOW_ADMIN` に変更していたとしても、
-- DEFAULT_ROLEが ACCOUNTADMIN / ORGADMIN / GLOBALORGADMIN / SECURITYADMIN だと
-- Openflowにアクセス拒否される
-- (`The role requiested has been explicity blocked ...` エラーが返る)
-- そのため、DEFAULT_ROLEを非特権のOPENFLOW_ADMINに変更して、セカンダリロールをALLにする
-- 変更後は一度サインアウト→再サインインしてからOpenflowを操作すること
ALTER USER <your_user> SET DEFAULT_ROLE = OPENFLOW_ADMIN;
ALTER USER <your_user> SET DEFAULT_SECONDARY_ROLES = ('ALL');
-- Openflowランタイムに紐づくロール
CREATE ROLE IF NOT EXISTS OPENFLOW_RUNTIME_ROLE_PGNLB;
-- ウェアハウス利用権限付与
GRANT USAGE, OPERATE ON WAREHOUSE OPENFLOW_WH TO ROLE OPENFLOW_RUNTIME_ROLE_PGNLB;
-- ランタイム作成時にこのランタイム用のロールを付与するために、Openflow管理ロールに権限付与
GRANT ROLE OPENFLOW_RUNTIME_ROLE_PGNLB TO ROLE OPENFLOW_ADMIN;
-- ネットワーク接続設定
CREATE DATABASE IF NOT EXISTS OPENFLOW_DB;
CREATE SCHEMA IF NOT EXISTS OPENFLOW_DB.NETWORKING;
CREATE OR REPLACE NETWORK RULE OPENFLOW_DB.NETWORKING.NLB_PG_EGRESS
MODE = EGRESS
TYPE = HOST_PORT
VALUE_LIST = ('<NLB_DNS>:5432');
CREATE OR REPLACE EXTERNAL ACCESS INTEGRATION OPENFLOW_PG_NLB_EAI
ALLOWED_NETWORK_RULES = (OPENFLOW_DB.NETWORKING.NLB_PG_EGRESS)
ENABLED = TRUE;
GRANT USAGE ON INTEGRATION OPENFLOW_PG_NLB_EAI TO ROLE OPENFLOW_RUNTIME_ROLE_PGNLB;
GRANT USAGE ON INTEGRATION OPENFLOW_PG_NLB_EAI TO ROLE OPENFLOW_ADMIN;
-- 連携先となるデータベース作成
CREATE DATABASE IF NOT EXISTS PG_OPENFLOW_NLB_DEST;
GRANT USAGE ON DATABASE PG_OPENFLOW_NLB_DEST TO ROLE OPENFLOW_RUNTIME_ROLE_PGNLB;
GRANT CREATE SCHEMA ON DATABASE PG_OPENFLOW_NLB_DEST TO ROLE OPENFLOW_RUNTIME_ROLE_PGNLB;
GRANT USAGE ON DATABASE PG_OPENFLOW_NLB_DEST TO ROLE OPENFLOW_ADMIN;
CREATE OR REPLACE EXTERNAL ACCESS INTEGRATION にて、外部ネットワークへの接続許可を行う必要があります。この設定内容は、後ほどOpenflowランタイムに反映させます。
7. Openflow Deploy設定
次に、Snowsight上からOpenflowデプロイメントの設定を行います。
作業ユーザを OPENFLOW_ADMIN ロールに切り替えて実施します。
Snowsight上の左メニューから「取り込み」>「Openflow」をクリックし、「Openflowを起動」をクリックします。
Openflowの画面が表示されたら、「Create a deployment」をクリックします。
「Prerequesites」の画面では「Next」をクリックし、「Deployment location」の画面では「Snowflake」を選択して、「Name」欄に適宜デプロイメント名を入力して「Next」をクリックします。

「Deployment configuration」では、今回は特に何も入力せずそのまま「Create deployment」をクリックします。
OpenflowデプロイメントのステータスがActiveになったら、次に進みましょう。
8. Openflow Runtime設定
次にOpenflow Runtimeを設定します。
Openflowの画面から「Create a runtime」をクリックします。
「Create runtime」の画面では以下の設定を行い、「Create」をクリックします。
- Deployment
- 作成したOpenflowデプロイメント名を選択
- Runtime Name
- 適宜入力する
- ここでは
PGNLBとした
- Node size
- Medium 以上を選択
- ここでは
Mediumを選択した
- Min nodes
- 1
- Max nodes
- 1
- PostgreSQLコネクタ要件により、マルチノードをサポートしていないため、単一ノードとなるようMin nodes/Max nodesとも
1を指定します- 参考: Openflow要件
- Execute-as role
OPENFLOW_RUNTIME_ROLE_PGNLBを選択


これでランタイムが作成されます。
作成には2分程度かかります。ランタイムのステータスが「Active」になるまで待ちましょう。
9. Openflowコネクタ導入
Openflow Runtimeが作成されたら、PostgreSQL用のコネクタを導入します。
Openflowのoverview画面内にある「Featured connectors」に「PostgreSQL」が表示されていれば、そのPostgreSQLパネルの「Install」をクリックします。

もし表示されていなければ、「View more connectors」リンクをクリックし、コネクタ一覧の中から「PostgreSQL」を検索して同様に「Install」をクリックします。

PostgreSQLコネクタインストール先となるOpenflow Runtimeを選択します。
ここでは先ほど作成したOpenflow Runtimeを選択しましょう。
「Add」ボタンをクリックすると、ランタイムにPostgreSQLコネクタが導入されます。

Snowflake 資格情報で認証し、ランタイムのアクセスを許可します。


成功すると、キャンバスにコネクタのプロセスグループが表示されます。

なお、ランタイムへのアクセス時に「The role requested has been explicitly blocked for use with this application」と表示されて、ログイン画面に遷移される場合、認証する作業ユーザの DEFAULT_ROLE の設定を変更する必要があります。
03-openflow-deployment-setup.sql の中で設定していますが、この部分の適用がまだの場合は再度設定をしてください。
ALTER USER <your_user> SET DEFAULT_ROLE = OPENFLOW_ADMIN;
ALTER USER <your_user> SET DEFAULT_SECONDARY_ROLES = ('ALL');
10. Openflowコネクタ設定
キャンバス上では、以下の設定を行います。
まず、PostgreSQLと書かれたボックスを右クリックして「Parameters」をクリックします。
Source / Destination / Ingestion の各値を以下のようにし、Source / Destination / Ingestion 毎に「Apply」をクリックして適用します。
「Parameters」クリック時には、Source / Destination / Ingestion パラメータ設定画面のいずれかがポップアップ表示されますので、適宜設定したら別のパラメータ設定に移りましょう。

Source (RDS for PostgreSQLへの接続設定)
| パラメータ | 値 |
|---|---|
| PostgreSQL Connection URL | jdbc:postgresql://<NLB_DNS>:5432/appdb?sslmode=require |
| PostgreSQL Username | openflow_repl |
| PostgreSQL Password | 02-publication-user.sql で設定したパスワード |
| Publication Name | snowflake_pub |
| PostgreSQL JDBC Driver | postgresql.org の JDBC jar をアップロードして指定 (「Reference asset」にもチェック) |
| Replication Slot Name | 空でよい(snowflake_connector_<random> のReplication Slotが自動作成される) |
JDBCのjarファイルアップロード方法は、こちらのDevelopersIO記事をご確認ください。
Destination (Snowflakeへの接続設定)
| パラメータ | 値 |
|---|---|
| Destination Database | PG_OPENFLOW_NLB_DEST |
| Destination Schema Pattern | 例: ${source.schema.name} (ソース側PostgreSQLの public スキーマを再現) |
| Snowflake Authentication Strategy | SNOWFLAKE_MANAGED |
| Snowflake Role | OPENFLOW_RUNTIME_ROLE_PGNLB |
| Snowflake Warehouse | OPENFLOW_WH |
Ingestion (取り込み設定)
| パラメータ | 値 |
|---|---|
| Included Table Names | public.customers,public.orders |
| Merge Task Schedule (CRON) | 例: 0 * * * * ? (1 分毎。お試し用。本番は要調整) |
11. Snowflake Egress IPをNLBのSecurity Groupに許可
フローを開始する前に、Snowsight上で以下のSQLを実行してSnowflake側のIPアドレス範囲を取得しておきます。
-- Snowsight で実行し、IPv4 CIDR を控える
SELECT SYSTEM$GET_SNOWFLAKE_EGRESS_IP_RANGES();
控えた CIDR を 1 行ずつに対して、NLB の Security Group のingress許可条件を追加します。
aws ec2 authorize-security-group-ingress \
--region "ap-northeast-1" \
--group-id "$NLB_SG" \
--ip-permissions \
"IpProtocol=tcp,FromPort=${PORT},ToPort=${PORT},IpRanges=[{CidrIp=${cidr},Description=snowflake-egress}]"
12. フローの開始
キャンバスのボックスでない部分を右クリックし、「Enable all Controller Services」をクリックします

続いて、キャンバスのボックスを右クリックし、「Start」をクリックします。

上記の順に実行するとコネクタが起動し、初回スナップショット → 増分(CDC)の順で取り込みを開始します。
動作確認
コネクタが起動したら、動作確認をします。
初回スナップショット
まずは初回スナップショットが成功している確認しましょう。
初回スナップショットが完了すると、PostgreSQLでpublication設定した対象スキーマ public と各テーブルがSnowflake上に反映されていることがわかります。
Snowsight上から、以下のSQLを実行してデータ件数を確認しましょう。
SELECT COUNT(*) FROM PG_OPENFLOW_NLB_DEST."public"."customers";
SELECT COUNT(*) FROM PG_OPENFLOW_NLB_DEST."public"."orders";
PostgreSQLに格納されているデータ件数と一致しているはずです。
CDC (差分転送)
次に、ソース側のPostgreSQLにてデータ更新を行い、その内容がSnowflakeへ反映されているか確認します。
PostgreSQLへ以下のSQLクエリを実行して、データの変更を行いました。
INSERT INTO public.orders (customer_id, amount, status) VALUES (3, 250.00, 'NEW');
UPDATE public.customers SET name = 'Alice T.', updated_at = now() WHERE id = 1;
DELETE FROM public.orders WHERE id = 2;
Merge スケジュール(1分毎)の後、Snowflake 側で反映を確認します。
SELECT * FROM PG_OPENFLOW_NLB_DEST."public"."customers" ORDER BY "id";
SELECT * FROM PG_OPENFLOW_NLB_DEST."public"."orders" ORDER BY "id";
まとめ
いかがでしたか。
Snowflake Openflow Connector for PostgreSQLを今回はSnowflakeのSPCS上に導入してみました。
事前の設定は、BYOC同様に思いのほか苦労しましたが、こちらも何とか動くところまでは確認できました。
今回はインターネット経由で接続していますが、SPCS上でOpenflow Connectorを動作する場合はPrivateLink経由での接続が安全です。
PrivateLink経由での接続の場合はSnowflakeのBusiness Criticalエディションが必要になるので、StandardエディションやEnterpriseエディションでSnowflakeを利用されている組織では、
今回のような接続パターンはSnowflake側のCIDRからのみにアクセス制限をかける等の対策は必要になるはずです。
この記事がお役に立てれば幸いです。








