【问题标题】:Authenticate with ECE ElasticSearch Sink from Apache Fink (Scala code)使用来自 Apache Fink 的 ECE ElasticSearch Sink 进行身份验证(Scala 代码)
【发布时间】: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 请求。其次,您可以在本地调试应用程序并在RestHighLevelClientconvertExistsResponse() 方法中设置断点,以查看您收到的响应。
  • 我正在尝试获取一个包含根 CA 证书的新证书,也许这有助于解决这个问题,因为那时我不必使用 trustfulSslContext,但这可能需要一天左右的时间至少。与 Flink 文档中提供的代码 sn-p 相比,对“new RestClientFactory”的使用有什么想法吗?我担心由于无法解决 Flink 文档中的代码问题,我自己编写了一些新内容,但现在出现了不同的问题。

标签: scala elasticsearch apache-flink flink-streaming


【解决方案1】:

很高兴报告此问题已解决。我在 Update 1 中发布的代码是正确的。对 ECE 的 ping 不起作用有两个原因:

  1. 证书需要包含完整的链,包括根 CA、中间 CA 和 ECE 证书。这有助于摆脱整个 trustfulSslContext 的东西。

  2. ECE 位于 ha-proxy 后面,代理将 HTTP 请求中的主机名映射到 ECE 中的实际部署集群名称。此映射逻辑没有考虑到 Java REST 高级客户端使用 org.apache.httphost 类,该类将主机名创建为 hostname:port_number 即使端口号为 443。由于 443 它没有找到映射因此,ECE 返回 404 错误而不是 200 ok(找到此错误的唯一方法是查看 ha-proxy 上的未加密数据包)。一旦修复了 ha-proxy 中的映射逻辑,就找到了映射并且 ping 成功了。

【讨论】:

    猜你喜欢
    • 2020-03-31
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-07-09
    相关资源
    最近更新 更多