我想通过 DoFn 向在 Dataflow 上运行的 Apache Beam Pipeline 发出 POST 请求。
为此,我创建了一个客户端,它实例化在 PoolingHttpClientConnectionManager 上配置的 HttpClosableClient。
但是,我为我处理的每个元素实例化一个客户端。
我如何设置一个由我的所有元素使用的持久客户端?
还有其他我应该使用的并行和高速 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)