All Products
Search
Document Center

Realtime Compute for Apache Flink:KMS credential management integration (Beta)

Last Updated:Sep 23, 2026

Realtime Compute for Apache Flink supports referencing credentials hosted in Alibaba Cloud Key Management Service (KMS) in jobs, replacing plaintext usernames and passwords (user-pwd) when accessing data sources. This provides encrypted credential hosting, avoids plaintext exposure, and enables maintenance-free credential management.

Overview

By integrating KMS credential management, you can reference KMS-hosted credentials (such as database usernames and passwords) in Flink jobs without exposing plaintext credentials in the job code. At runtime, Flink assumes a RAM role to obtain a temporary STS token, calls KMS to retrieve the actual credential value, and then connects to the target data source.

Key benefits:

  • Sensitive information protection: No plaintext user-pwd is exposed in jobs. Credentials are hosted encrypted in KMS.

  • Maintenance-free credentials: With KMS credential rotation, Flink automatically fetches the latest credential value after updates. Job code does not need to be modified.

Notes

Incorrectly deleting credentials in KMS, or incorrectly deleting or modifying the permissions and trust policy of a RAM role in RAM, may cause Flink jobs to fail because they cannot obtain KMS credentials. According to the Service Level Agreement for Realtime Compute for Apache Flink, unavailability caused by incorrect resource usage by customers is not covered by the SLA.

Limits

  • Only the Postgres CDC Connector (VVR 11.8 and later) and the Kafka Connector (VVR 11.9 and later) support accessing cloud resources with KMS credentials.

  • There is a limit on the number of referenced RAM roles for accessing KMS. For details, see Use a RAM role to access cloud resources.

Permission preparation

  • Step 1 (create a KMS credential): The operator must have KMS credential management permissions.

  • Step 2 and Step 3: The operator must be granted the AliyunRAMFullAccess policy, or separately granted RAM role management and permission management permissions. For more information, see Create a custom policy.

  • Step 4 and Step 5 (reference a role in Flink and develop jobs): The operator must have editor permissions or higher on the Flink namespace, or be granted fine-grained permissions to reference RAM roles, manage files, and develop jobs. For details, see Grant permissions to access the development console.

Step 1: Create a KMS credential

KMS provides multiple types of credentials. The following steps use an RDS credential as an example:

  1. Log on to the Key Management Service console. Select the region on the top menu bar, and in the left-side navigation pane, choose Resources > Secrets.

  2. On the Self-Managed Secrets tab, in the Secret Type area, click Database Secret.

  3. Above the list on the right, click Create Single Secret under Create Database Secret, complete the configuration, and click OK.

For more information, see Get started with Secrets Manager and Manage and use ApsaraDB RDS secrets.

Step 2: Create and configure a RAM role

This role is the identity credential that the Flink service uses to access KMS. Flink assumes the role to obtain temporary credentials and calls KMS as the role to retrieve credential values. For details about the credential-free access to KMS, see the appendix of Use a RAM role to access cloud resources.

1. Create a RAM role

  1. Log on to the RAM console - Roles page and click Create Role. Keep the default options and click OK.

  2. Enter a recognizable role name, for example, FlinkRoleForKMSRead, and click OK.

2. Add a trust policy

After the role is created, on the role details page, click the Trust Policy tab and click Edit Trust Policy. Add "Service": ["stream.aliyuncs.com"] to Principal to trust Flink to assume the role. For more information, see Modify the trust policy of a RAM role.

{
  "Statement": [
    {
      "Action": "sts:AssumeRole",
      "Effect": "Allow",
      "Principal": {
        "Service": [
          "stream.aliyuncs.com"
        ]
      }
    }
  ],
  "Version": "1"
}

3. Grant KMS access permissions to the RAM role

After Flink assumes the role, it calls KMS as the role. Grant the required KMS permissions to the role's permission policy in advance.

On the role details page, click the Permissions tab and click Grant Permission. Add the AliyunKMSSecretUserAccess policy (which grants permission to retrieve credentials in KMS). If credentials are encrypted with a customer master key (CMK), also grant the kms:Decrypt permission by adding the AliyunKMSCryptoUserAccess policy.

To control permissions at the key and credential level, create a custom policy that grants GetSecretValue and Decrypt permissions and restricts the resource to specified resources. For details, see KMS custom policy reference and Create a custom policy.

{
  "Version": "1",
  "Statement": [
    {
      "Effect": "Allow",
      "Action": "kms:GetSecretValue",
      "Resource": "acs:kms:<region>:<account-id>:secret/<secret-name>"
    },
    {
      "Effect": "Allow",
      "Action": "kms:Decrypt",
      "Resource": "acs:kms:<region>:<account-id>:key/<key-id>"
    }
  ]
}

Step 3: Grant permissions to a RAM user

Create a permission policy and grant it to the RAM user. This allows the user to pass the RAM role created in Step 2 to the Flink service for assumption.

Note

If the RAM user in Step 4 is already associated with the AliyunRAMFullAccess or AliyunStreamFullAccess policy, the user already has the PassRole permission. You can skip this step.

Create a permission policy

Log on to the RAM console - Policies page and click Create Policy. Switch to Script mode and edit the policy content.

{
  "Version": "1",
  "Statement": [
    {
      "Effect": "Allow",
      "Action": [
        "ram:ListRoles"
      ],
      "Resource": "*"
    },
    {
      "Effect": "Allow",
      "Action": "ram:PassRole",
      "Resource": "acs:ram::<account-id>:role/<role-name>",
      "Condition": {
        "StringEquals": {
          "acs:Service": "stream.aliyuncs.com"
        }
      }
    }
  ]
}

Permission

Description

ram:ListRoles

Allows the user to view the role list under the account, so that the role can be selected in the Flink console.

ram:PassRole

Allows the user to pass the specified role to the Flink cloud service. Replace <account-id> and <role-name> in the policy with actual values:

  • <account-id>: The UID of the Alibaba Cloud account.

  • <role-name>: The name of the RAM role created in Step 2.

Add a Condition to restrict passing the role to stream.aliyuncs.com only.

Click OK to save the policy, for example, as FlinkPassRolePolicy. For more information, see Create a custom policy.

Grant permissions to the RAM user

  1. Log on to the RAM console - Users page. Find the RAM user to authorize and click Add Permissions in the Actions column.

  2. In the Add Permissions panel, search for and select the policy created in Step 3 (for example, FlinkPassRolePolicy), then click Confirm. For more information, see Grant permissions to a RAM user.

Step 4: Reference the RAM role in Flink

Reference the RAM role created in Step 2 to the Flink namespace so that jobs can use the role to access KMS.

  1. Log on to the Realtime Compute console and enter the target workspace.

  2. In the left-side navigation pane, choose Security > Security.

  3. Click the RAM Roles tab.

  4. Click Reference RAM Role.

  5. In the dialog box, select the RAM role created in Step 2 (for example, FlinkRoleForKMSRead). Fuzzy search by role name is supported.

    Note

    The role list contains regular service roles under the Alibaba Cloud account that owns the current namespace. The operator must belong to the same Alibaba Cloud account as the current namespace and have the ram:ListRoles permission to view the roles.

  6. Click Test. The system automatically checks the following two items:

    • PassRole permission check: The current user has the ram:PassRole permission on the selected RAM role.

    • Trust policy check: The selected RAM role trusts the current Flink namespace to assume it.

  7. After the checks pass, click OK.

For more information, see Use a RAM role to access cloud resources.

Step 5: Use KMS credentials in jobs

After the previous steps are complete, specify credentials using the secret:// reference format in the job code and configure KMS connection parameters. Flink automatically retrieves credential values at runtime to access the data source.

Postgres CDC Connector example

SQL job: The following example uses a Postgres CDC source table and hosts the database username and password with a KMS credential. A database credential named my-db-secret has been created in KMS, with SecretData {"AccountName":"alice","AccountPassword":"xxx"}.

CREATE TEMPORARY TABLE postgrescdc_source (
    id INT NOT NULL,
    name STRING,
    description STRING,
    weight DECIMAL(10,3)
) WITH (
  'connector' = 'postgres-cdc',
  'hostname' = '<host name>',
  'port' = '<port>',
  'username' = 'secret://kms.my-db-secret.AccountName',
  'password' = 'secret://kms.my-db-secret.AccountPassword',
  'database-name' = '<database name>',
  'schema-name' = '<schema name>',
  'table-name' = '<table name>',
  'slot.name' = '<slot name>',
  'decoding.plugin.name' = 'pgoutput',
  'kms.endpoint' = '<your-kms-instance-vpc-endpoint>',
  'kms.ca' = 'oss://<bucket>/kms-ca.pem',
  'kms.akless.assume-role.role-name' = 'FlinkRoleForKMSRead'
);

For the parameter descriptions, prerequisites, and type mapping of the Postgres CDC connector, see Postgres CDC connector.

Kafka Connector example

Kafka credentials are referenced in the SASL JAAS configuration. Replace the username and password values in properties.sasl.jaas.config with secret:// references. Other Kafka parameters stay the same. A credential named my-kafka-secret has been created in KMS, with SecretData {"AccountName":"<SASL username>","AccountPassword":"<SASL password>"}.

SQL job: The following example uses a Kafka source table. Configure a result table in the same way.

CREATE TEMPORARY TABLE kafka_source (
  order_id BIGINT,
  user_id BIGINT,
  amount INT,
  payload STRING
) WITH (
  'connector' = 'kafka',
  'topic' = '<topic name>',
  'properties.bootstrap.servers' = '<broker 1>,<broker 2>,<broker 3>',
  'properties.group.id' = '<group id>',
  'scan.startup.mode' = 'earliest-offset',
  'format' = 'csv',
  'properties.security.protocol' = 'SASL_PLAINTEXT',
  'properties.sasl.mechanism' = 'PLAIN',
  /* Reference KMS credentials in username and password */
  'properties.sasl.jaas.config' = 'org.apache.flink.kafka.shaded.org.apache.kafka.common.security.plain.PlainLoginModule required username="secret://kms.my-kafka-secret.AccountName" password="secret://kms.my-kafka-secret.AccountPassword";',
  'kms.endpoint' = 'kms.cn-beijing.aliyuncs.com',
  'kms.akless.assume-role.role-name' = 'FlinkRoleForKMSRead'
);

Data ingestion YAML job: The source and the sink specify KMS parameters separately. The source uses a shared gateway and the sink uses a dedicated gateway, so this example covers both.

source:
  type: kafka
  name: Kafka Source
  topic: <topic name>
  value.format: json
  scan.startup.mode: earliest-offset
  properties.bootstrap.servers: <broker 1>,<broker 2>,<broker 3>
  properties.group.id: <group id>
  properties.security.protocol: SASL_PLAINTEXT
  properties.sasl.mechanism: PLAIN
  properties.sasl.jaas.config: org.apache.flink.kafka.shaded.org.apache.kafka.common.security.plain.PlainLoginModule required username="secret://kms.my-kafka-secret.AccountName" password="secret://kms.my-kafka-secret.AccountPassword";
  # Shared gateway: the endpoint is fixed per region, and kms.ca is not required
  kms.endpoint: kms.cn-beijing.aliyuncs.com
  kms.akless.assume-role.role-name: FlinkRoleForKMSRead

sink:
  type: kafka
  name: Kafka Sink
  properties.bootstrap.servers: <broker 1>,<broker 2>,<broker 3>
  # ApsaraMQ for Kafka does not support idempotent writes. Disable idempotence on VVR 8.0.0 and later
  properties.enable.idempotence: false
  properties.security.protocol: SASL_PLAINTEXT
  properties.sasl.mechanism: PLAIN
  properties.sasl.jaas.config: org.apache.flink.kafka.shaded.org.apache.kafka.common.security.plain.PlainLoginModule required username="secret://kms.my-kafka-secret.AccountName" password="secret://kms.my-kafka-secret.AccountPassword";
  # Dedicated gateway: kms.ca is also required
  kms.endpoint: <your-kms-instance-vpc-endpoint>
  kms.ca: oss://flink-fullymanaged-<workspace-id>/artifacts/namespaces/<namespace>/PrivateKmsCA_kst-******.pem
  kms.akless.assume-role.role-name: FlinkRoleForKMSRead
Note
  • In YAML, the value of properties.sasl.jaas.config contains double quotes. Write them as they are, without escaping.

  • Get the endpoint addresses and port for SASL access from the ApsaraMQ for Kafka console.

For the complete parameters and security settings of the Kafka connector, see Kafka SQL connector and Kafka YAML connector.

Parameter description

Credential reference format

The KMS credential reference format is secret://<provider>.<secret-name>.<key>.

Field

Description

<provider>

The credential source type identifier. Currently, only kms is supported.

<secret-name>

The credential name in KMS, for example, my-db-secret.

<key>

The field name in the credential's SecretData JSON, for example, AccountName or AccountPassword.

Example: secret://kms.my-db-secret.AccountPassword

KMS parameters

Parameter

Type

Required

Description

kms.endpoint

String

Yes

The KMS gateway address (dedicated gateway or shared gateway).

  • Dedicated gateway:

    • On the Instance Management page of KMS, click the Software Key Management or Hardware Key Management tab, and then select the target KMS instance.

    • Click the instance ID to go to the details page and view the VPC address of the instance, which is the endpoint. For more information, see Service endpoints.

  • Shared gateway: the endpoint is fixed by region and can be obtained from Service endpoints.

kms.ca

String

Required for dedicated gateways

The OSS path of the CA certificate.

  1. Download the CA certificate from the details page of the KMS instance. The downloaded file is named PrivateKmsCA_kst-******.pem by default.

  2. Upload the certificate by using Flink Artifacts, and then copy the OSS path in the format oss://bucket/path.

kms.akless.assume-role.role-name

String

No

The name of the RAM role used to access KMS. The role must be referenced in advance under Security > RAM Roles.