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
AliyunRAMFullAccesspolicy, 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:
-
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.
-
On the Self-Managed Secrets tab, in the Secret Type area, click Database Secret.
-
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
-
Log on to the RAM console - Roles page and click Create Role. Keep the default options and click OK.
-
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.
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 |
|
|
Allows the user to view the role list under the account, so that the role can be selected in the Flink console. |
|
|
Allows the user to pass the specified role to the Flink cloud service. Replace
Add a |
Click OK to save the policy, for example, as FlinkPassRolePolicy. For more information, see Create a custom policy.
Grant permissions to the RAM user
-
Log on to the RAM console - Users page. Find the RAM user to authorize and click Add Permissions in the Actions column.
-
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.
-
Log on to the Realtime Compute console and enter the target workspace.
-
In the left-side navigation pane, choose .
-
Click the RAM Roles tab.
-
Click Reference RAM Role.
-
In the dialog box, select the RAM role created in Step 2 (for example,
FlinkRoleForKMSRead). Fuzzy search by role name is supported.NoteThe 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:ListRolespermission to view the roles. -
Click Test. The system automatically checks the following two items:
-
PassRole permission check: The current user has the
ram:PassRolepermission on the selected RAM role. -
Trust policy check: The selected RAM role trusts the current Flink namespace to assume it.
-
-
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
-
In YAML, the value of
properties.sasl.jaas.configcontains 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 |
|
|
The credential source type identifier. Currently, only |
|
|
The credential name in KMS, for example, |
|
|
The field name in the credential's SecretData JSON, for example, |
Example: secret://kms.my-db-secret.AccountPassword
KMS parameters
|
Parameter |
Type |
Required |
Description |
|
|
String |
Yes |
The KMS gateway address (dedicated gateway or shared gateway).
|
|
|
String |
Required for dedicated gateways |
The OSS path of the CA certificate.
|
|
|
String |
No |
The name of the RAM role used to access KMS. The role must be referenced in advance under Security > RAM Roles. |