【问题标题】:Error using SpannerIO in apache beam在 Apache Beam 中使用 SpannerIO 时出错
【发布时间】:2018-03-22 21:12:50
【问题描述】:

这个问题是this one 的后续问题。 我正在尝试使用 apache 梁从谷歌扳手表中读取数据(然后进行一些数据处理)。我使用 java SDK 编写了以下最小示例:

package com.google.cloud.dataflow.examples;
import java.io.IOException;
import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.PipelineResult;
import org.apache.beam.sdk.io.gcp.spanner.SpannerIO;
import org.apache.beam.sdk.options.PipelineOptions;
import org.apache.beam.sdk.options.PipelineOptionsFactory;
import org.apache.beam.sdk.values.PCollection;
import com.google.cloud.spanner.Struct;

public class backup {

  public static void main(String[] args) throws IOException {
    PipelineOptions options = PipelineOptionsFactory.create();

    Pipeline p = Pipeline.create(options);
    PCollection<Struct> rows = p.apply(
            SpannerIO.read()
                .withInstanceId("my_instance")
                .withDatabaseId("my_db")
                .withQuery("SELECT t.table_name FROM information_schema.tables AS t")
                );
    
    PipelineResult result = p.run();
    try {
      result.waitUntilFinish();
    } catch (Exception exc) {
      result.cancel();
    }
  }
}

当我尝试使用 DirectRunner 执行代码时,我得到 以下错误信息:

org.apache.beam.runners.direct.repackaged.com.google.common.util.concurrent.UncheckedExecutionException:

org.apache.beam.sdk.util.UserCodeException: java.lang.NoClassDefFoundError:无法初始化类 com.google.cloud.spanner.spi.v1.SpannerErrorInterceptor

[...] 原因: org.apache.beam.sdk.util.UserCodeException: java.lang.NoClassDefFoundError:无法初始化类 com.google.cloud.spanner.spi.v1.SpannerErrorInterceptor

[...] 原因: java.lang.NoClassDefFoundError:无法初始化类 com.google.cloud.spanner.spi.v1.SpannerErrorInterceptor

或者,使用 DataflowRunner:

org.apache.beam.runners.direct.repackaged.com.google.common.util.concurrent.UncheckedExecutionException: org.apache.beam.sdk.util.UserCodeException: java.lang.NoSuchFieldError: internal_static_google_rpc_LocalizedMessage_fieldAccessorTable

[...] 引起:org.apache.beam.sdk.util.UserCodeException: java.lang.NoSuchFieldError: internal_static_google_rpc_LocalizedMessage_fieldAccessorTable

[...] 引起:java.lang.NoSuchFieldError: internal_static_google_rpc_LocalizedMessage_fieldAccessorTable

在这两种情况下,错误消息都相当神秘,我无法从谷歌搜索中找到任何关于导致错误的明确想法。我也找不到任何使用 SpannerIO 模块的示例脚本。

这个错误是由于我的代码中有明显错误,还是由于谷歌云工具安装错误?

【问题讨论】:

标签: java google-cloud-dataflow apache-beam google-cloud-spanner


【解决方案1】:

此问题很可能是由此处描述的依赖项兼容性问题引起的:BEAM-2837。以下是 JIRA 问题中的一个 cmets 中描述的快速解决方法:

<dependency>
    <groupId>com.google.api.grpc</groupId>
    <artifactId>grpc-google-common-protos</artifactId>
    <version>0.1.9</version>
</dependency>

<dependency>
    <groupId>org.apache.beam</groupId>
    <artifactId>beam-sdks-java-io-google-cloud-platform</artifactId>
    <version>${beam.version}</version>
    <exclusions>
        <exclusion>
            <groupId>com.google.api.grpc</groupId>
            <artifactId>grpc-google-common-protos</artifactId>
        </exclusion>
    </exclusions>
</dependency>

明确定义所需的com.google.api.grpc 依赖项并从org.apache.beam 中排除版本。

【讨论】:

  • 谢谢!不过老实说,我最终还是继续使用 python SDK,并制作了一个自定义 ParDo 来读取/写入扳手。
【解决方案2】:

您需要指定 ProjectID:

    SpannerIO.read()
            .withProjectId("my_project")
            .withInstanceId("my_instance")
            .withDatabaseId("my_db")

您需要为您的 Spanner 项目设置凭据。由于 SpannerIO 的 API 不允许您设置任何自定义凭据,因此您必须使用环境变量 GOOGLE_APPLICATION_CREDENTIALS 设置全局应用程序凭据。

您还可以使用 JDBC 读取(和写入)Cloud Spanner。阅读是这样完成的:

        PCollection<KV<String, Long>> words = p2.apply(JdbcIO.<KV<String, Long>> read()
            .withDataSourceConfiguration(JdbcIO.DataSourceConfiguration.create("nl.topicus.jdbc.CloudSpannerDriver",
                    "jdbc:cloudspanner://localhost;Project=my-project-id;Instance=instance-id;Database=database;PvtKeyPath=C:\\Users\\MyUserName\\Documents\\CloudSpannerKeys\\cloudspanner-key.json"))
            .withQuery("SELECT t.table_name FROM information_schema.tables AS t").withCoder(KvCoder.of(StringUtf8Coder.of(), BigEndianLongCoder.of()))
            .withRowMapper(new JdbcIO.RowMapper<KV<String, Long>>()
            {
                private static final long serialVersionUID = 1L;

                @Override
                public KV<String, Long> mapRow(ResultSet resultSet) throws Exception
                {
                    return KV.of(resultSet.getString(1), resultSet.getLong(2));
                }
            }));

此方法还允许您通过设置 PvtKeyPath 来使用自定义凭据。您还可以使用 JDBC 写入 Google Cloud Spanner。看看这里的例子:http://www.googlecloudspanner.com/2017/10/google-cloud-spanner-with-apache-beam.html

【讨论】:

  • 我确实忘记了“projectID”行,尽管添加它并不能解决错误。事实上,我正在使用 Eclipse 谷歌云工具插件,并登录我的谷歌帐户。所以这应该照顾凭据?我可能得试试 JDBC 版本。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-06-23
  • 1970-01-01
  • 2019-01-25
  • 2018-09-22
相关资源
最近更新 更多