相关疑难解决方法(0)

DoFn 中的 HTTP 客户端

我想通过 DoFn 向在 Dataflow 上运行的 Apache Beam Pipeline 发出 POST 请求。

为此,我创建了一个客户端,它实例化在 PoolingHttpClientConnectionManager 上配置的 HttpClosableClient。

但是,我为我处理的每个元素实例化一个客户端。

我如何设置一个由我的所有元素使用的持久客户端?

还有其他我应该使用的并行和高速 HTTP 请求类吗?

http google-cloud-dataflow apache-beam apache-beam-io

5
推荐指数
1
解决办法
2454
查看次数

如何使用 Apache Beam (Java) 进行异步 Http 调用?

输入PCollection是http请求,它是一个有界数据集。我想在 ParDo 中进行异步 http 调用(Java),解析响应并将结果放入输出 PCollection 中。我的代码如下。获取异常如下。

我不明白原因。需要指导....

java.util.concurrent.CompletionException: java.lang.IllegalStateException: Can't add element ValueInGlobalWindow{value=streaming.mapserver.backfill.EnrichedPoint@2c59e, pane=PaneInfo.NO_FIRING} to committed bundle in PCollection Call Map Server With Rate Throttle/ParMultiDo(ProcessRequests).output [PCollection]
Run Code Online (Sandbox Code Playgroud)

代码:

public class ProcessRequestsFn extends DoFn<PreparedRequest,EnrichedPoint> {
    private static AsyncHttpClient _HttpClientAsync;
    private static ExecutorService _ExecutorService;

static{

    AsyncHttpClientConfig cg = config()
            .setKeepAlive(true)
            .setDisableHttpsEndpointIdentificationAlgorithm(true)
            .setUseInsecureTrustManager(true)
            .addRequestFilter(new RateLimitedThrottleRequestFilter(100,1000))
            .build();

    _HttpClientAsync = asyncHttpClient(cg);

    _ExecutorService = Executors.newCachedThreadPool();

}


@DoFn.ProcessElement
public void processElement(ProcessContext c) {

    PreparedRequest request = c.element();

    if(request == null)
        return;

    _HttpClientAsync.prepareGet((request.getRequest()))
            .execute()
            .toCompletableFuture()
            .thenApply(response -> …
Run Code Online (Sandbox Code Playgroud)

asynchttpclient apache-beam

2
推荐指数
1
解决办法
4366
查看次数