【问题标题】:Google Cloud Data Fusion -- building pipeline from REST API endpoint sourceGoogle Cloud Data Fusion——从 REST API 端点源构建管道
【发布时间】:2020-01-31 00:52:16
【问题描述】:

尝试构建管道以从第 3 方 REST API 端点数据源中读取数据。

我正在使用 Hub 中的 HTTP(1.2.0 版)插件。

响应请求网址为:https://api.example.io/v2/somedata?return_count=false

响应正文示例:

{
  "paging": {
    "token": "12456789",
    "next": "https://api.example.io/v2/somedata?return_count=false&__paging_token=123456789"
  },
  "data": [
    {
      "cID": "aerrfaerrf",
      "first": true,
      "_id": "aerfaerrfaerrf",
      "action": "aerrfaerrf",
      "time": "1970-10-09T14:48:29+0000",
      "email": "example@aol.com"
    },
    {...}
  ]
}

日志中的主要错误是:

java.lang.NullPointerException: null
    at io.cdap.plugin.http.source.common.pagination.BaseHttpPaginationIterator.getNextPage(BaseHttpPaginationIterator.java:118) ~[1580429892615-0/:na]
    at io.cdap.plugin.http.source.common.pagination.BaseHttpPaginationIterator.ensurePageIterable(BaseHttpPaginationIterator.java:161) ~[1580429892615-0/:na]
    at io.cdap.plugin.http.source.common.pagination.BaseHttpPaginationIterator.hasNext(BaseHttpPaginationIterator.java:203) ~[1580429892615-0/:na]
    at io.cdap.plugin.http.source.batch.HttpRecordReader.nextKeyValue(HttpRecordReader.java:60) ~[1580429892615-0/:na]
    at io.cdap.cdap.etl.batch.preview.LimitingRecordReader.nextKeyValue(LimitingRecordReader.java:51) ~[cdap-etl-core-6.1.1.jar:na]
    at org.apache.spark.rdd.NewHadoopRDD$$anon$1.hasNext(NewHadoopRDD.scala:214) ~[spark-core_2.11-2.3.3.jar:2.3.3]
    at org.apache.spark.InterruptibleIterator.hasNext(InterruptibleIterator.scala:37) ~[spark-core_2.11-2.3.3.jar:2.3.3]
    at scala.collection.Iterator$$anon$12.hasNext(Iterator.scala:439) ~[scala-library-2.11.8.jar:na]
    at scala.collection.Iterator$$anon$12.hasNext(Iterator.scala:439) ~[scala-library-2.11.8.jar:na]
    at scala.collection.Iterator$$anon$12.hasNext(Iterator.scala:439) ~[scala-library-2.11.8.jar:na]
    at org.apache.spark.internal.io.SparkHadoopWriter$$anonfun$4.apply(SparkHadoopWriter.scala:128) ~[spark-core_2.11-2.3.3.jar:2.3.3]
    at org.apache.spark.internal.io.SparkHadoopWriter$$anonfun$4.apply(SparkHadoopWriter.scala:127) ~[spark-core_2.11-2.3.3.jar:2.3.3]
    at org.apache.spark.util.Utils$.tryWithSafeFinallyAndFailureCallbacks(Utils.scala:1415) ~[spark-core_2.11-2.3.3.jar:2.3.3]
    at org.apache.spark.internal.io.SparkHadoopWriter$.org$apache$spark$internal$io$SparkHadoopWriter$$executeTask(SparkHadoopWriter.scala:139) [spark-core_2.11-2.3.3.jar:2.3.3]
    at org.apache.spark.internal.io.SparkHadoopWriter$$anonfun$3.apply(SparkHadoopWriter.scala:83) [spark-core_2.11-2.3.3.jar:2.3.3]
    at org.apache.spark.internal.io.SparkHadoopWriter$$anonfun$3.apply(SparkHadoopWriter.scala:78) [spark-core_2.11-2.3.3.jar:2.3.3]
    at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:87) [spark-core_2.11-2.3.3.jar:2.3.3]
    at org.apache.spark.scheduler.Task.run(Task.scala:109) [spark-core_2.11-2.3.3.jar:2.3.3]
    at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:345) [spark-core_2.11-2.3.3.jar:2.3.3]
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) [na:1.8.0_232]
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) [na:1.8.0_232]
    at java.lang.Thread.run(Thread.java:748) [na:1.8.0_232]

可能的问题

在尝试解决此问题一段时间后,我认为问题可能出在

分页

  • Data Fusion HTTP插件有很多处理分页的方法
    • 根据上面的响应正文,分页类型的最佳选择似乎是Link in Response Body
    • 对于所需的 Next Page JSON/XML Field Path 参数,我尝试了$.paging.nextpaging/next。两者都不起作用。
    • 我已验证/paging/next 中的链接在 Chrome 中打开时有效

认证

  • 当只是尝试在 Chrome 中查看响应 URL 时,会弹出一个提示,询问用户名和密码
    • 只需输入用户名的 API 密钥即可在 Chrome 中通过此提示
    • 要在 Data Fusion HTTP 插件中执行此操作,API 密钥用于 基本身份验证部分中的用户名

在数据源是 REST API 的 Google Cloud Data Fusion 中,任何人都成功地创建了管道?

【问题讨论】:

    标签: rest pagination endpoint google-cloud-data-fusion


    【解决方案1】:

    回答

    在数据源是 REST API 的 Google Cloud Data Fusion 中,任何人都成功地创建了管道?

    这不是实现这一目标的最佳方法,最好的方法是将数据 Service APIs Overview 提取到发布/订阅,然后使用发布/订阅作为管道的源,这将为您提供一个简单可靠的暂存位置用于处理、存储和分析的数据,请参阅 pub/sub API 的文档。为了将它与 Dataflow 结合使用,要遵循的步骤在此处的官方文档中Using Pub/Sub with Dataflow

    【讨论】:

    • 感谢您的回复。您能否详细说明为什么 Data Fusion 不适合此应用程序?仅仅是因为 Pub/Sub 提供了可靠的暂存位置吗?看起来 Data Fusion 可以简单地将 Cloud Storage 用作可靠的暂存位置。
    • 问题不是数据融合,而是根据您的描述,我了解到您希望直接从 REST 端点摄取...您当然可以将数据融合与 pub/sub 一起使用...如果您查看本教程,它会给您一个更好的想法codelabs.developers.google.com/codelabs/real-time-csv-cdf-bq/…
    • 是的,这是正确的:这里的真正问题是如何从 REST 端点摄取数据。但是,我仍然对 Pub/Sub 如何帮助我向我的数据源 API 发送 GET 请求并在转移到 Data Fusion 或 Dataflow 之前暂存响应数据感到困惑。您在 OP 中提供的服务 API 链接是关于 Pub/Sub API,但我看不到它涉及从谷歌自己的 API 以外的 API 端点摄取数据的地方
    • @Korean_Of_the_Mountain 你应该看看 Apache NiFi。
    【解决方案2】:

    我认为您的问题在于您收到的数据格式。例外:

    java.lang.NullPointerException: null
    

    当您没有指定正确的输出架构时发生(我相信在这种情况下没有架构)

    解决方案 1

    要解决这个问题,请尝试将 HTTP 数据融合插件配置为:

    • 接收格式:文本。
    • 输出架构:名称:用户类型:字符串

    这应该可以从 API 以字符串格式获取响应。完成后,使用 JSONParser 将字符串转换为类似对象的表格。

    解决方案 2

    将 HTTP 数据融合插件配置为:

    • 接收格式:json
    • JSON/XML 结果路径:数据
    • JSON/XML 字段映射:包括您提供的字段(见附图)。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2021-08-11
      • 1970-01-01
      • 2022-09-26
      • 1970-01-01
      • 1970-01-01
      • 2019-09-15
      • 1970-01-01
      相关资源
      最近更新 更多