【发布时间】:2019-05-25 23:33:02
【问题描述】:
使用 Flink 文档中提供的示例时出现编译器错误。 Flink 文档提供了示例 Scala 代码,用于在与 Elasticsearch 通信时设置 REST 客户端工厂参数,https://ci.apache.org/projects/flink/flink-docs-stable/dev/connectors/elasticsearch.html。 尝试此代码时,我在 IntelliJ 中遇到编译器错误,提示“无法解析符号 restClientBuilder”。
我发现以下 SO 完全是我的问题,除了它是在 Java 中并且我在 Scala 中执行此操作。 Apache Flink (v1.6.0) authenticate Elasticsearch Sink (v6.4)
我尝试将上述 SO 中提供的解决方案代码复制粘贴到 IntelliJ 中,自动转换的代码也有编译器错误。
// provide a RestClientFactory for custom configuration on the internally created REST client
// i only show the setMaxRetryTimeoutMillis for illustration purposes, the actual code will use HTTP cutom callback
esSinkBuilder.setRestClientFactory(
restClientBuilder -> {
restClientBuilder.setMaxRetryTimeoutMillis(10)
}
)
然后我尝试了(由 IntelliJ 自动生成 Java 到 Scala 代码)
// provide a RestClientFactory for custom configuration on the internally created REST client// provide a RestClientFactory for custom configuration on the internally created REST client
import org.apache.http.auth.AuthScope
import org.apache.http.auth.UsernamePasswordCredentials
import org.apache.http.client.CredentialsProvider
import org.apache.http.impl.client.BasicCredentialsProvider
import org.apache.http.impl.nio.client.HttpAsyncClientBuilder
import org.elasticsearch.client.RestClientBuilder
// provide a RestClientFactory for custom configuration on the internally created REST client// provide a RestClientFactory for custom configuration on the internally created REST client
esSinkBuilder.setRestClientFactory((restClientBuilder) => {
def foo(restClientBuilder) = restClientBuilder.setHttpClientConfigCallback(new RestClientBuilder.HttpClientConfigCallback() {
override def customizeHttpClient(httpClientBuilder: HttpAsyncClientBuilder): HttpAsyncClientBuilder = { // elasticsearch username and password
val credentialsProvider = new BasicCredentialsProvider
credentialsProvider.setCredentials(AuthScope.ANY, new UsernamePasswordCredentials(es_user, es_password))
httpClientBuilder.setDefaultCredentialsProvider(credentialsProvider)
}
})
foo(restClientBuilder)
})
原始代码 sn -p 产生错误“无法解析 RestClientFactory”,然后 Java 到 Scala 显示其他几个错误。
所以基本上我需要找到Apache Flink (v1.6.0) authenticate Elasticsearch Sink (v6.4)中描述的解决方案的Scala版本
更新 1:在 IntelliJ 的帮助下,我取得了一些进展。下面的代码可以编译运行,但是还有一个问题。
esSinkBuilder.setRestClientFactory(
new RestClientFactory {
override def configureRestClientBuilder(restClientBuilder: RestClientBuilder): Unit = {
restClientBuilder.setHttpClientConfigCallback(new RestClientBuilder.HttpClientConfigCallback() {
override def customizeHttpClient(httpClientBuilder: HttpAsyncClientBuilder): HttpAsyncClientBuilder = {
// elasticsearch username and password
val credentialsProvider = new BasicCredentialsProvider
credentialsProvider.setCredentials(AuthScope.ANY, new UsernamePasswordCredentials(es_user, es_password))
httpClientBuilder.setDefaultCredentialsProvider(credentialsProvider)
httpClientBuilder.setSSLContext(trustfulSslContext)
}
})
}
}
问题是我不确定我是否应该做一个新的 RestClientFactory 对象。发生的情况是应用程序连接到 elasticsearch 集群,但随后发现 SSL CERT 无效,所以我不得不放置 trustfullSslContext(如此处所述https://gist.github.com/iRevive/4a3c7cb96374da5da80d4538f3da17cb),这让我解决了 SSL 问题,但现在是 ES REST客户端执行 ping 测试,但 ping 失败,它抛出异常并且应用程序关闭。我怀疑 ping 失败是因为 SSL 错误,也许它没有使用我设置的 trustfulSslContext 作为新 RestClientFactory 的一部分,这让我怀疑我不应该做新的,应该有一个简单的方法来更新现有的 RestclientFactory 对象,基本上这一切都是由于我缺乏 Scala 知识而发生的。
【问题讨论】:
-
失败时会收到什么错误?
-
当 Update 1 中显示的代码失败时,我收到“Java.lang.RuntimeException:没有可访问的 Elasticsearch 节点!”错误。我猜这是因为github.com/apache/flink/blob/master/flink-connectors/… 中的内置 ping 测试仍然使用旧的 REST 客户端对象,该对象不使用我通过 setRestClientFactory 设置的 trustfulSslContext。
-
好吧,据我所知。
ping()基本上向给定地址发出带有给定参数的HEAD请求,并简单地返回布尔值,说明响应代码是否为 200。您可以在这里做两件事。首先,您可以尝试向您的弹性搜索所在的地址发出HEAD请求。其次,您可以在本地调试应用程序并在RestHighLevelClient的convertExistsResponse()方法中设置断点,以查看您收到的响应。 -
我正在尝试获取一个包含根 CA 证书的新证书,也许这有助于解决这个问题,因为那时我不必使用 trustfulSslContext,但这可能需要一天左右的时间至少。与 Flink 文档中提供的代码 sn-p 相比,对“new RestClientFactory”的使用有什么想法吗?我担心由于无法解决 Flink 文档中的代码问题,我自己编写了一些新内容,但现在出现了不同的问题。
标签: scala elasticsearch apache-flink flink-streaming