如何修改EventHub AAD认证项目,用证书替代密钥实现PySpark流消费?
问题:Azure Synapse中PySpark Structured Streaming消费EventHub改用证书认证
我希望在Azure Synapse中使用PySpark的Structured Streaming消费EventHub消息,有人推荐了相关项目,现在想修改项目,用证书认证替代原本的密钥认证。我尝试修改下面的ServicePrincipalCredentialsAuth类,使用ClientCredentialFactory.createFromCertificate方法,但没能成功。
package net.alexott.demos.eventhubs_aad; import com.microsoft.aad.msal4j.ClientCredentialFactory; import com.microsoft.aad.msal4j.ConfidentialClientApplication; import com.microsoft.aad.msal4j.IClientCredential; import java.net.MalformedURLException; import scala.collection.immutable.Map; public class ServicePrincipalCredentialsAuth extends ServicePrincipalAuthBase { private final String clientSecret; private static final String AAD_CLIENT_SECRET_KEY = "aad_client_secret"; public ServicePrincipalCredentialsAuth(Map<String, String> params) { super(params); clientSecret = params.get(AAD_CLIENT_SECRET_KEY).get(); } @Override ConfidentialClientApplication getClient() throws MalformedURLException { IClientCredential credential = ClientCredentialFactory.createFromSecret(this.clientSecret); return ConfidentialClientApplication.builder(this.clientId, credential) .authority(this.authEndpoint) .build(); } }
解决方案:修改为证书认证实现
直接替换原类的逻辑,改为通过证书加载凭证,以下是完整修改后的代码:
package net.alexott.demos.eventhubs_aad; import com.microsoft.aad.msal4j.ClientCredentialFactory; import com.microsoft.aad.msal4j.ConfidentialClientApplication; import com.microsoft.aad.msal4j.IClientCredential; import java.io.FileInputStream; import java.io.IOException; import java.net.MalformedURLException; import java.security.KeyStore; import java.security.KeyStoreException; import java.security.NoSuchAlgorithmException; import java.security.cert.CertificateException; import scala.collection.immutable.Map; public class ServicePrincipalCertificateAuth extends ServicePrincipalAuthBase { private final String certPath; private final String certPassword; // 定义证书相关参数键名 private static final String AAD_CERT_PATH_KEY = "aad_cert_path"; private static final String AAD_CERT_PASSWORD_KEY = "aad_cert_password"; public ServicePrincipalCertificateAuth(Map<String, String> params) { super(params); // 从参数中获取证书路径和密码 certPath = params.get(AAD_CERT_PATH_KEY).get(); certPassword = params.get(AAD_CERT_PASSWORD_KEY).get(); } @Override ConfidentialClientApplication getClient() throws MalformedURLException, KeyStoreException, IOException, NoSuchAlgorithmException, CertificateException { // 加载PKCS12格式的证书文件(PFX证书通用此格式) KeyStore keyStore = KeyStore.getInstance("PKCS12"); try (FileInputStream fis = new FileInputStream(certPath)) { keyStore.load(fis, certPassword.toCharArray()); } // 通过证书创建认证凭证 IClientCredential credential = ClientCredentialFactory.createFromCertificate(keyStore, certPassword.toCharArray()); return ConfidentialClientApplication.builder(this.clientId, credential) .authority(this.authEndpoint) .build(); } }
关键注意事项
- 替换项目中原有的
ServicePrincipalCredentialsAuth引用为新类ServicePrincipalCertificateAuth - 确保证书文件(PFX格式)能被Synapse集群访问:可上传至ADLS Gen2的挂载路径,或放置在集群本地存储目录
- 调整参数配置:将原有的
aad_client_secret替换为aad_cert_path(证书文件路径)和aad_cert_password(证书密码),保留aad_client_id、aad_tenant_id等原有参数 - 验证Azure AD配置:确保服务主体已上传对应证书,且证书未过期、权限配置正确
- 异常处理:新增了证书加载相关的异常类型,需要在调用处处理或向上抛出
内容的提问来源于stack exchange,提问作者s528060
相关产品推荐
相关产品推荐

