Apache Beam调用ElasticsearchIO读取AWS OpenSearch报403错误如何解决
问题根因
403报错的核心原因是Apache Beam内置的
ElasticsearchIO默认不会对发往AWS OpenSearch集群的请求做AWS SigV4签名,你在PipelineOptions中配置的AWS凭证不会被该IO算子的HTTP客户端自动读取使用,OpenSearch集群收到未签名的请求会判定为匿名用户访问,触发权限拦截。
解决方案
你需要自定义Elasticsearch客户端的请求拦截器,给所有发往OpenSearch的请求加上SigV4签名,再将自定义客户端配置传入ElasticsearchIO中,具体操作步骤如下:
1. 引入必要依赖(以Maven为例)
<dependency> <groupId>com.amazonaws</groupId> <artifactId>aws-java-sdk-core</artifactId> <version>1.12.XXX</version> <!-- 替换为和你项目适配的版本 --> </dependency> <dependency> <groupId>org.elasticsearch.client</groupId> <artifactId>elasticsearch-rest-client</artifactId> <version>7.10.2</version> <!-- 要和你OpenSearch兼容的客户端版本一致 --> </dependency>
2. 实现AWS SigV4签名请求拦截器
import com.amazonaws.auth.AWSCredentialsProvider; import com.amazonaws.http.HttpMethodName; import com.amazonaws.auth.Signer; import com.amazonaws.auth.DefaultRequestSigner; import com.amazonaws.DefaultRequest; import com.amazonaws.Request; import com.amazonaws.util.IOUtils; import org.apache.http.HttpRequest; import org.apache.http.HttpRequestInterceptor; import org.apache.http.protocol.HttpContext; import java.io.ByteArrayInputStream; import java.net.URI; import java.util.HashMap; import java.util.Map; public class AwsSigV4RequestInterceptor implements HttpRequestInterceptor { private static final String SERVICE_NAME = "es"; private final String region; private final AWSCredentialsProvider credentialsProvider; private final Signer signer; public AwsSigV4RequestInterceptor(String region, AWSCredentialsProvider credentialsProvider) { this.region = region; this.credentialsProvider = credentialsProvider; this.signer = DefaultRequestSigner.builder().serviceName(SERVICE_NAME).region(region).build(); } @Override public void process(HttpRequest request, HttpContext context) throws Exception{ String method = request.getRequestLine().getMethod(); URI uri = URI.create(request.getRequestLine().getUri()); // 构造AWS签名请求 Request<?> awsRequest = new DefaultRequest<>(SERVICE_NAME); awsRequest.setHttpMethod(HttpMethodName.fromValue(method)); awsRequest.setEndpoint(URI.create("https://" + request.getFirstHeader("Host").getValue())); awsRequest.setResourcePath(uri.getRawPath()); // 处理请求参数 if (uri.getRawQuery() != null) { Map<String, String> params = new HashMap<>(); for (String param : uri.getRawQuery().split("&")) { String[] kv = param.split("=", 2); params.put(kv[0], kv.length > 1 ? kv[1] : ""); } awsRequest.setParameters(params); } // 处理请求体 if (request instanceof org.apache.http.HttpEntityEnclosingRequest) { org.apache.http.HttpEntity entity = ((org.apache.http.HttpEntityEnclosingRequest) request).getEntity(); if (entity != null) { awsRequest.setContent(new ByteArrayInputStream(IOUtils.toByteArray(entity.getContent()))); } } // 签名并把签名头加到原请求中 signer.sign(awsRequest, credentialsProvider.getCredentials()); for (Map.Entry<String, String> header : awsRequest.getHeaders().entrySet()) { request.setHeader(header.getKey(), header.getValue()); } } }
3. 修改ElasticsearchIO配置,传入带签名拦截器的客户端
ElasticSearchOptions options = PipelineOptionsFactory.fromArgs(args).withValidation() .as(ElasticSearchOptions.class); Pipeline pipeline = Pipeline.create(options); // 初始化凭证和拦截器 AWSCredentialsProvider awsCredentialsProvider = new AWSStaticCredentialsProvider( new BasicAWSCredentials(options.getAwsAccessKey(), options.getAwsSecretKey()) ); AwsSigV4RequestInterceptor signInterceptor = new AwsSigV4RequestInterceptor("us-east-1", awsCredentialsProvider); // 配置带拦截器的RestClient ElasticsearchIO.ConnectionConfiguration esConfig = ElasticsearchIO.ConnectionConfiguration.create( <hostName>, options.getElasticSearchIndex(), options.getElasticSearchType() ) .withClientConfigurationFn(restClientBuilder -> { restClientBuilder.setHttpClientConfigCallback(httpClientBuilder -> httpClientBuilder.addInterceptorLast(signInterceptor) ); return restClientBuilder; }); PCollection<String> output = pipeline.apply( ElasticsearchIO.read() .withConnectionConfiguration(esConfig) .withQuery("field_name:some-value") );
注意事项
- 请确认你的AWS访问密钥对应的IAM身份,已经被授予目标OpenSearch集群的
es:ESHttpGet权限,IAM策略配置无误 - 签名用的区域必须和你OpenSearch集群实际部署的区域完全一致
- 如果你的OpenSearch集群配置了IP访问白名单,需要将GCP Beam工作节点的出口IP段加入白名单,或者在IAM策略中移除IP限制条件
内容的提问来源于stack exchange,提问作者zerodark
相关产品推荐
相关产品推荐

