Snowflake Openflow Connector for PostgreSQL を SPCS で動かし、RDS for PostgreSQLに接続してみた

Snowflake Openflow Connector for PostgreSQL を SPCS で動かし、RDS for PostgreSQLに接続してみた

SnowflakeのOpenflow Connector for PostgreSQLを使って、RDS for PostgreSQLからSnowflakeへCDCデータ転送を試してみました。前回はAWS上にデプロイしてデプロイして動作するBYOCを試しましたが、今回はSPCSで試してみました。
2026.08.18

データ事業本部の笠原です。

先日SnowflakeのOpenflow Connector for PostgreSQLを使って、RDS for PostgreSQLのデータをCDC (Change Data Capture) でSnowflakeに転送できるか、試しました。

この時は、OpenflowのデプロイモデルのうちBYOCを採用し、自身のAWSアカウント上のEKS上でOpenflowを動かしてみました。

今回はもう一つのデプロイモデルであるSPCSを採用して、Snowflake上でOpenflowを動かして、RDS for PostgreSQLに接続できるか試してみました。

今回の構成

今回の構成は以下の通りです。

architecture_nlb

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
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; で確認しましょう。
  • 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」をクリックします。

openflow_deployment01

「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 を指定します
  • Execute-as role
    • OPENFLOW_RUNTIME_ROLE_PGNLB を選択

openflow_runtime01

openflow_runtime02

これでランタイムが作成されます。
作成には2分程度かかります。ランタイムのステータスが「Active」になるまで待ちましょう。

9. Openflowコネクタ導入

Openflow Runtimeが作成されたら、PostgreSQL用のコネクタを導入します。

Openflowのoverview画面内にある「Featured connectors」に「PostgreSQL」が表示されていれば、そのPostgreSQLパネルの「Install」をクリックします。

openflow_connector_postgres00

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

openflow_connector_postgres01

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

openflow_connector_postgres01

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

openflow_access01

openflow_access02

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

openflow_connector_postgres04

なお、ランタイムへのアクセス時に「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 パラメータ設定画面のいずれかがポップアップ表示されますので、適宜設定したら別のパラメータ設定に移りましょう。

openflow_connector01

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」をクリックします

openflow_connector02

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

openflow_connector03

上記の順に実行するとコネクタが起動し、初回スナップショット → 増分(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からのみにアクセス制限をかける等の対策は必要になるはずです。

この記事がお役に立てれば幸いです。


Snowflake World Tour Tokyo 2026に参加しませんか?

Snowflakeの国内最大級イベント「Snowflake World Tour Tokyo」が2026年9月10日(水)・11日(木)にグランドプリンスホテル新高輪にて開催されます。
最新のAI・データ活用事例やライブデモを体感できる無料イベントです。

Snowflake World Tour Tokyoイベントに参加する


Snowflakeの導入支援はクラスメソッドに!

クラスメソッドでは Snowflake の導入を支援しております。
製品の詳細や支援の内容についてお気軽にお問い合わせください。

Snowflakeの詳細を見る

この記事をシェアする

AWSのお困り事はクラスメソッドへ

関連記事