You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何配置Firehose客户端并解决X509TrustManager执行错误

问题描述

使用AWS Firehose数据流的putRecord API,已用Kotlin实现接收Firehose客户端和投递流名称的RecordPublisherImpl类,通过Dagger注入Firehose客户端,并在Lambda组件的execute方法中调用该类的发布方法。运行时抛出以下错误:

software.amazon.awssdk.core.exception.SdkClientException: Unable to execute HTTP request: No X509TrustManager implementation available

怀疑是客户端配置不当(缺少凭证),需解决两个问题:

  1. 如何通过ARN或access_key、secret_key配置Firehose客户端
  2. 修复上述X509TrustManager相关错误

相关代码

RecordPublisherImpl 实现代码

class RecordPublisherImpl @Inject constructor(
    private val deliveryStreamName: String,
    private val firehoseClient: FirehoseClient,
) : EventRecordPublisher {
   // 其他业务函数
   fun publishRecord(*params*){
       // 预处理逻辑 - 将参数转换为recordString
       publishRecordToDataFirehose(recordString)
   }

   private fun publishRecordToDataFirehose(
        recordString: String
    ): Boolean {
        val eventRecordBytes = recordString.toByteArray()
        val eventRecordSdkBytes = SdkBytes.fromByteArray(eventRecordBytes)
        val eventRecordAsDataFirehoseRecord = Record.builder().data(eventRecordSdkBytes).build()
        val eventRecordAsPutRecordRequest = PutRecordRequest.builder().deliveryStreamName(deliveryStreamName)
            .record(eventRecordAsDataFirehoseRecord).build()
        try {
            val putRecordResponse = firehoseClient.putRecord(eventRecordAsPutRecordRequest)
            log.info("向$deliveryStreamName写入记录成功,响应:$putRecordResponse")
        } catch (e: RuntimeException) {
            throw IllegalStateException("向流$deliveryStreamName写入记录时发生异常", e)
        }
        return true
    }
}

Dagger Module 配置代码

@Module
class AccessorModule {

    @Provides
    @Singleton
    fun provideAwsCredentialsProvider(): AwsCredentialsProvider {
        return DefaultCredentialsProvider.create()
    }

    @Provides
    @Singleton
    fun provideFirehoseClient(
        awsCredentialsProvider: AwsCredentialsProvider
    ): FirehoseClient{
        return FirehoseClient
            .builder()
            .region(Region.of("region"))
            .credentialsProvider(awsCredentialsProvider)
            .build()
    }

    @Provides
    @Singleton
    @Named("DeliveryStreamName")
    fun providesDeliveryStreamName(): String = "streamName"

    @Provides
    @Singleton
    fun provideEventRecordPublisherImpl(@Named("DeliveryStreamName") deliveryStreamName: String, firehoseClient: FirehoseClient): RecordPublisherImpl{
        return RecordPublisherImpl(deliveryStreamName, firehoseClient)
    }
}

Lambda 组件代码

class GenerateActivity
@Inject constructor(
    private val recordPublisher: RecordPublisherImpl
) {
    override fun execute(input: GenerateInput) {
        log.debug("收到输入:{}", input)

        log.debug("调用recordPublisher...")
        val putRecordResponse = recordPublisher.publishRecord(
            "123",
            "String1",
            "SomeString",
            "SomethingElse"
        )
        log.debug("收到响应:{}", putRecordResponse)
    }
}

解决方案

一、修复X509TrustManager错误

这个错误并非直接由凭证缺失导致,而是运行环境缺少SSL证书管理依赖。AWS SDK for Java v2需要SSL组件建立HTTPS连接,可按以下方式解决:

  • JVM环境:确保项目依赖包含完整的aws-sdk-core,或显式添加SSL相关依赖(比如Gradle中添加org.apache.httpcomponents:httpclient)
  • Lambda自定义运行时:若使用自定义运行时,需确保运行时包含完整JRE环境(默认Lambda运行时已自带,无需额外配置)
  • 显式配置SSL上下文:在Firehose客户端构建时手动指定SSL上下文,示例代码如下:
import software.amazon.awssdk.http.apache.ApacheHttpClient
import javax.net.ssl.SSLContext

// 修改provideFirehoseClient方法
fun provideFirehoseClient(awsCredentialsProvider: AwsCredentialsProvider): FirehoseClient{
    val sslContext = SSLContext.getDefault()
    val httpClient = ApacheHttpClient.builder()
        .sslContext(sslContext)
        .build()
        
    return FirehoseClient
        .builder()
        .region(Region.of("your-region"))
        .credentialsProvider(awsCredentialsProvider)
        .httpClient(httpClient)
        .build()
}

二、配置Firehose客户端凭证

方式1:使用access_key和secret_key直接配置

修改provideAwsCredentialsProvider方法,用StaticCredentialsProvider注入凭证(生产环境禁止硬编码,建议通过环境变量注入):

@Provides
@Singleton
fun provideAwsCredentialsProvider(): AwsCredentialsProvider {
    val accessKey = System.getenv("AWS_ACCESS_KEY_ID")
    val secretKey = System.getenv("AWS_SECRET_ACCESS_KEY")
    val credentials = AwsBasicCredentials.create(accessKey, secretKey)
    return StaticCredentialsProvider.create(credentials)
}

方式2:使用IAM角色ARN(跨账号/托管服务场景)

如果在Lambda、EC2等AWS托管服务中运行,优先使用IAM角色。若需指定特定角色ARN(比如跨账号访问),可使用AssumeRoleCredentialsProvider:

import software.amazon.awssdk.auth.credentials.AssumeRoleCredentialsProvider
import software.amazon.awssdk.services.sts.StsClient

@Provides
@Singleton
fun provideAwsCredentialsProvider(): AwsCredentialsProvider {
    val stsClient = StsClient.builder().region(Region.of("your-region")).build()
    return AssumeRoleCredentialsProvider.builder()
        .roleArn("arn:aws:iam::123456789012:role/your-firehose-role")
        .roleSessionName("firehose-session")
        .stsClient(stsClient)
        .build()
}

需确保该IAM角色拥有firehose:PutRecord权限,且当前环境具备调用STS AssumeRole的权限。

方式3:使用默认凭证链

原代码中的DefaultCredentialsProvider.create()会自动按以下顺序查找凭证,无需额外配置:

  1. 系统环境变量(AWS_ACCESS_KEY_ID和AWS_SECRET_ACCESS_KEY)
  2. Java系统属性(aws.accessKeyId和aws.secretKey)
  3. 本地凭证文件(默认路径~/.aws/credentials)
  4. AWS托管服务的IAM角色凭证

三、其他排查点

  • 确认投递流名称和区域配置与实际Firehose流一致
  • 检查IAM权限:确保使用的凭证拥有firehose:PutRecord权限,资源范围包含目标投递流的ARN
  • Lambda环境下,确认执行角色已关联正确的权限策略

内容的提问来源于stack exchange,提问作者KG_1

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.25 09:15:27