Apache Beam 是一种统一的开源模型,用于 定义批次数据和流式数据的并行处理流水线。本文档介绍了如何在 Apache Beam 流水线中使用 SpannerIO 连接器,以从 Spanner Omni 数据库读取数据或向其中写入数据。
准备工作
如需将 SpannerIO 连接到 Spanner Omni,请确保满足以下要求:
在 Spanner Omni 环境中初始化数据库。
如果您使用加密,请确保使用兼容版本的 Apache Beam:
- 对于 TLS 加密,版本 2.69.0 或更高版本。
- 对于双向 TLS (mTLS) 加密,版本 2.75.0 或更高版本。
为您的环境设置身份验证凭据。
配置 SpannerIO 以连接到 Spanner Omni
如需将 SpannerIO 连接到 Spanner Omni,请使用数据库详细信息和连接参数配置 SpannerConfig。
如需配置连接,请选择以下连接模式之一:
使用纯文本通信进行连接
如需建立纯文本连接,请指定 Spanner Omni 端点,使用 withExperimentalHost() 方法启用实验性主机支持,并使用 withUsingPlainTextChannel() 方法配置流水线。
以下示例展示了如何配置纯文本连接:
SpannerConfig spannerConfig =
SpannerConfig.create()
.withDatabaseId("DATABASE_ID")
// Define the Spanner Omni endpoint
.withExperimentalHost("http://ENDPOINT")
// Use a plain-text connection
.withUsingPlainTextChannel(true);
替换以下内容:
DATABASE_ID:Spanner Omni 数据库的 ID,例如test-db。ENDPOINT:Spanner Omni 实例的端点,例如localhost:15000。
使用加密进行连接
如需保护数据库流量并确保 Apache Beam 和 Spanner Omni 之间的安全通信,您可以使用 TLS 或 mTLS 加密进行连接。 使用加密有助于确保凭据和数据的机密性。
使用 TLS 加密
如需使用 TLS 加密保护 Apache Beam 和 Spanner Omni 之间的数据库流量,您无需在 SpannerConfig 中指定凭据属性。相反,请使用 Spanner Omni CA 证书配置 Java 信任库,然后配置 SpannerConfig 以使用安全的 TLS 端点。
第 1 步:配置 Java 信任库
如需确保通信安全,您必须将 Spanner Omni 生成的 CA 证书导入 Java 信任库。使用以下任一选项:
默认 Java 信任库
运行以下命令,将 Spanner Omni 生成的 CA 证书添加到标准 Java 信任库:
sudo keytool -import -trustcacerts \
-file ~/.spanner/certs/ca.crt \
-alias spanner-ca \
-keystore $JAVA_HOME/lib/security/cacerts
自定义信任库
如需确保流水线仍可以连接到使用标准证书授权机构 (CA) 的其他数据库或服务,请创建自定义信任库:
复制现有 Java 信任库,以创建自定义信任库:
cp $JAVA_HOME/lib/security/cacerts PATH_TO_CUSTOM_CA_CERTIFICATE将 CA 证书导入自定义信任库:
keytool -import -trustcacerts \ -file ~/.spanner/certs/ca.crt \ -alias spanner-ca \ -keystore PATH_TO_CUSTOM_CA_CERTIFICATE运行流水线时传递自定义 CA 证书存储区:
java -Djavax.net.ssl.trustStore=PATH_TO_CUSTOM_CA_CERTIFICATE \ -Djavax.net.ssl.trustStorePassword=changeit \ -jar PIPELINE_NAME.jar
替换以下内容:
PATH_TO_CUSTOM_CA_CERTIFICATE:自定义 CA 证书存储区的路径。PIPELINE_NAME:Apache Beam 流水线的名称。
第 2 步:配置 SpannerConfig
如需配置 SpannerConfig 以使用安全的 TLS 连接,请将以下代码添加到流水线:
SpannerConfig spannerConfig =
SpannerConfig.create()
.withDatabaseId("DATABASE_ID")
// Define the secure Spanner Omni endpoint
.withExperimentalHost("https://ENDPOINT");
替换以下内容:
DATABASE_ID:Spanner Omni 数据库的 ID,例如test-db。ENDPOINT:Spanner Omni 实例的端点,例如localhost:15000。
使用 mTLS 加密
如需使用 Apache Beam 建立双向 TLS (mTLS) 连接,您必须使用 CA 证书配置 Java 信任库,生成客户端私钥或将其转换为 PKCS#8 格式,然后使用客户端证书和密钥路径配置 SpannerConfig。
第 1 步:配置 Java 信任库
使用 Spanner Omni CA 证书配置 Java 信任库,如本文档前面 第 1 步:配置 Java 信任库中所述。
第 2 步:转换或生成客户端私钥
如需使用 mTLS 进行连接,请确保客户端私钥采用 PKCS#8 格式。使用以下任一选项:
openssl
如需将 Spanner Omni 生成的客户端密钥转换为符合 Java 的格式,请运行以下命令:
openssl pkcs8 -topk8 \
-in ~/.spanner/certs/client.key \
-out ~/.spanner/certs/java-client.key \
-nocrypt
Spanner Omni CLI
使用 Spanner Omni CLI 和 --generate-pkcs8-key 标志创建客户端证书时,直接以 PKCS#8 格式生成密钥。
如需以 PKCS#8 格式生成客户端证书和客户端私钥,请运行以下命令:
spanner certificates create-client CLIENT_NAME \
--ca-certificate-directory=PATH_TO_CA_CERTIFICATES \
--ca-private-key-directory=PATH_TO_PRIVATE_KEYS \
--output-directory=PATH_TO_CERTIFICATES \
--generate-pkcs8-key
替换以下内容:
CLIENT_NAME:要为其生成证书和私钥的客户端的名称。PATH_TO_CA_CERTIFICATES:包含 CA 证书的目录的路径。PATH_TO_PRIVATE_KEYS:包含 CA 私钥的目录的路径。PATH_TO_CERTIFICATES:保存客户端证书和私钥的目录的路径。
第 3 步:配置 SpannerConfig
使用客户端证书和客户端私钥配置流水线代码中的 SpannerConfig:
SpannerConfig spannerConfig =
SpannerConfig.create()
.withDatabaseId("DATABASE_ID")
// Define the secure Spanner Omni endpoint
.withExperimentalHost("https://ENDPOINT")
// Specify the paths to the client certificate and private key
.withClientCert(
"PATH_TO_CLIENT_CERT",
"PATH_TO_CLIENT_CERT_KEY");
替换以下内容:
DATABASE_ID:Spanner Omni 数据库的 ID,例如test-db。ENDPOINT:Spanner Omni 实例的端点,例如localhost:15000。PATH_TO_CLIENT_CERT:客户端证书文件的路径。PATH_TO_CLIENT_CERT_KEY:客户端私钥文件的路径。
配置身份验证令牌
不建议客户端使用身份验证令牌,因为 Spanner Omni 生成的令牌会过期,并且需要使用 Spanner Omni CLI 手动续订。如需将身份验证令牌与 Spanner Omni 端点的 TLS 或 mTLS 设置搭配使用,请将 SPANNER_EXPERIMENTAL_HOST_AUTH_TOKEN 环境变量设置为 Spanner Omni CLI 生成的身份验证令牌的值。对于不需要凭据的连接,请将此变量保留为未设置状态。