ALBのルーティングを応用してレガシーWEBシステムにモダンなサブシステムを導入する 3. ALB/Kinesis編

Agaroot IT Partners(AITP)のtomoです。

前々回の記事から紹介している、ALBを用いたWEBシステムにECS on Fargateによって構成されたサブシステムを追加するCDKコードの実装方法についてです。今までの記事は下記をご参照ください。

  1. 概要編
  2. ECR/タスク定義編

今回はALB/Kinesisを中心に触れていきます。

ALBの設定

ALBに関する実装を見ていきましょう。ターゲットグループを作成し、既存リスナーへ追加していきます。
主なソースはこちらになります。
https://github.com/aitp-tomo/alb-subsystem/blob/main/lib/wrapper/AlbWrapper.ts

  private readonly getListener = (): void => {
    this.listener = elbv2.ApplicationListener.fromApplicationListenerAttributes(
      this.scope,
      `${this.appId}-listener`,
      {
        listenerArn: `arn:aws:elasticloadbalancing:${this.region}:${this.account}:listener/app/${this.listenerId}`,
        securityGroup: this.securityGroup,
      }
    );
  };

  private readonly addTarget = (): void => {
    const targetGroupId = `${this.appId}-tg`;
    this.targetGroup = new elbv2.ApplicationTargetGroup(
      this.scope,
      targetGroupId,
      {
        healthCheck: {
          enabled: true,
          healthyHttpCodes: "200",
          path: this.healthcheckPath,
          port: "80",
          protocol: elbv2.Protocol.HTTP,
          timeout: cdk.Duration.seconds(this.healthCheckTimeout),
          interval: cdk.Duration.seconds(this.healthCheckInterval),
        },
        port: 80,
        protocol: elbv2.ApplicationProtocol.HTTP,
        protocolVersion: elbv2.ApplicationProtocolVersion.HTTP1,
        targets: [this.ecsWrapper.service],
        vpc: this.vpcWrapper.vpc,
        targetGroupName: targetGroupId,
      }
    );

    this.listener.addTargetGroups(`${this.appId}-listener-tg`, {
      targetGroups: [this.targetGroup],
      conditions: [
        elbv2.ListenerCondition.pathPatterns(this.listenerRulePathPatterns),
      ],
      priority: this.listenerRulePriority,
    });
  };

ALBで用いるターゲットグループの作成はelbv2.ApplicationTargetGroupによって行います。
https://docs.aws.amazon.com/cdk/api/v2/docs/aws-cdk-lib.aws_elasticloadbalancingv2.ApplicationTargetGroup.html
targetsに作成したECSサービスを含み、protocol及びportでコンテナが転送されたリクエストを受け付けられるように設定しましょう。また、healthCheckでヘルスチェックの指定も行えます。

作成したターゲットグループを別途elbv2.ApplicationListener.fromApplicationListenerAttributesで取得していた既存リスナーに追加していきます。
その際にconditionsの方でリスナールールの条件を指定します。今回は特定のパスの場合にサブシステム用の新規ターゲットグループへ転送したいので、elbv2.ListenerCondition.pathPatternsでパスを指定します。ホストヘッダなど他の条件も指定できますので、必要に応じて下記レファレンスを参照してください。
https://docs.aws.amazon.com/cdk/api/v2/docs/aws-cdk-lib.aws_elasticloadbalancingv2.ListenerCondition.html

Kinesisなどを用いたログの長期保存

単にサービスを動かすだけならばここまででも良いのですが、比較的高額なCloudWatch Logsには一時的に保存するだけにとどめ、別途S3バケットにログを転送して長期保存できるようにします。

具体的なアーキテクチャ図が以下のようになります。

ECSサービスにより出力されたCloudWatch LogsのログをKinesis Data Streams(以下よりKDS)及びKinesis Data Firehose(KDF)を用いてS3バケットに転送していきます。KDS及びKDFの概要については下記記事などをご参照ください。
https://dev.classmethod.jp/articles/difference-between-kinesis-streams-and-kinesis-firehose/

ECSとCloudWatch Logsとを結びつけるには、まずロググループを作成して、そのロググループをタスク定義のログドライバーというもので指定する必要があります。
https://docs.aws.amazon.com/ja_jp/AmazonECS/latest/developerguide/using_awslogs.html

これをCDKコードで行うにはlogs.LogGroupecs.AwsLogDriverを用いましょう。
https://github.com/aitp-tomo/alb-subsystem/blob/main/lib/wrapper/TaskDefinitionWrapper.ts

  private readonly createLogGroup = (): void => {
    const logGroupId = `${this.appId}-task-definition-log-group`;
    this.logGroup = new logs.LogGroup(this.scope, logGroupId, {
      logGroupName: logGroupId,
      retention: logs.RetentionDays.TWO_WEEKS,
    });
  };

  private readonly createLogDriver = (): void => {
    this.logDriver = new ecs.AwsLogDriver({
      streamPrefix: this.containerName,
      mode: ecs.AwsLogDriverMode.NON_BLOCKING,
      logGroup: this.logGroup,
    });
  };

  private readonly addContainer = (): void => {
    const image = ecs.EcrImage.fromDockerImageAsset(
      this.ecrWrapper.dockerImageAsset
    );
    this.taskDefinition.addContainer(this.containerName, {
      image,
      containerName: this.containerName,
      cpu: this.cpu,
      memoryLimitMiB: this.memoryLimitMiB,
      portMappings: [
        {
          containerPort: 80,
        },
      ],
      logging: this.logDriver,
    });

ログの保存先となるS3バケットの作成については以前公開した記事が参考になるかと思いますので、そちらをご参照ください。

CloudWatch LogsのロググループからKDSの方でログを収集するにはサブスクリプションフィルターにKDSのストリームを転送先に指定する必要があります。

これをCDKコードで実現する場合にはkinesis.Streamでストリームを作成し、logs.LogGroupaddSubscriptionFilterメソッドでそのストリームを転送先に指定します。
https://github.com/aitp-tomo/alb-subsystem/blob/main/lib/wrapper/KinesisWrapper.ts

  private readonly createSourceStream = (): void => {
    const sourceStreamId = `${this.appId}-source-stream`;
    this.scourceStream = new kinesis.Stream(this.scope, sourceStreamId, {
      streamName: sourceStreamId,
      streamMode: kinesis.StreamMode.ON_DEMAND,
    });
  };

  private readonly addSubscriptionFilter = (): void => {
    const filterId = `${this.appId}-filter`;
    this.taskDefinitionWrapper.logGroup.addSubscriptionFilter(filterId, {
      filterName: filterId,
      destination: new logsDestinations.KinesisDestination(this.scourceStream),
      filterPattern: logs.FilterPattern.allEvents(),
    });
  };

そしてKDFでKDSからのデータをS3バケットに保存するには、KDFで用いるIAMロールにS3バケット操作の権限を付与して、KDFのストリームを設定します。

CDKにおいてKDFのL2コンストラクトは存在しないので、レファレンスを参考にしながらL1コンストラクトのkinesisfirehose.CfnDeliveryStreamを実装していってください。
https://docs.aws.amazon.com/cdk/api/v2/docs/aws-cdk-lib.aws_kinesisfirehose.CfnDeliveryStream.html

  private readonly createRole = (): void => {
    const roleId = `${this.appId}-delivery-stream-role`;
    this.role = new iam.Role(this.scope, roleId, {
      assumedBy: new iam.ServicePrincipal("firehose.amazonaws.com"),
      roleName: roleId,
    });
    this.grantedByStream = this.scourceStream.grant(
      this.role,
      "kinesis:DescribeStream",
      "kinesis:GetShardIterator",
      "kinesis:GetRecords"
    );
    this.role.addToPolicy(
      new iam.PolicyStatement({
        actions: [
          "s3:AbortMultipartUpload",
          "s3:GetBucketLocation",
          "s3:GetObject",
          "s3:ListBucket",
          "s3:ListBucketMultipartUploads",
          "s3:PutObject",
        ],
        effect: iam.Effect.ALLOW,
        resources: [
          this.s3Wrapper.bucket.bucketArn,
          `${this.s3Wrapper.bucket.bucketArn}/*`,
        ],
      })
    );
    this.role.addToPolicy(
      new iam.PolicyStatement({
        actions: ["logs:PutLogEvents"],
        effect: iam.Effect.ALLOW,
        resources: [this.taskDefinitionWrapper.logGroup.logGroupArn],
      })
    );
  };

  private readonly createDeliveryStream = (): void => {
    const deliveryStreamId = `${this.appId}-delivery-stream`;
    const deliveryStream = new kinesisfirehose.CfnDeliveryStream(
      this.scope,
      deliveryStreamId,
      {
        deliveryStreamName: deliveryStreamId,
        deliveryStreamType: "KinesisStreamAsSource",
        kinesisStreamSourceConfiguration: {
          kinesisStreamArn: this.scourceStream.streamArn,
          roleArn: this.role.roleArn,
        },
        s3DestinationConfiguration: {
          bucketArn: this.s3Wrapper.bucket.bucketArn,
          roleArn: this.role.roleArn,
        },
      }
    );
    this.grantedByStream.applyBefore(deliveryStream);
  };

これによりS3バケットにログが長期保存できるようになります。

また、Kinesisを利用しているので必要に応じてデータ変換を行ってからS3バケットに保存したり、S3バケット以外の転送先(例えばOpenSearchなど)の追加もできるでしょう。

まとめ

今回は3本の記事を通してALBを用いた既存システムにサブシステムを導入する方法を取り上げました。

レガシーとなってしまっているシステムに追加機能を実装したい時、今回の記事を参考に保守性の高いモダンな技術スタックを用いたサブシステムを構築してみては如何でしょうか。

関連するタグ