【问题标题】:TupleTag Tag <taginfo> corresponds to a non-singleton resultTupleTag Tag <taginfo> 对应一个非单例结果
【发布时间】:2017-09-22 02:25:09
【问题描述】:

在 Google Cloud Dataflow 中,我的加入失败并显示“TupleTag 标记对应于非单例结果” 从错误堆栈看来,这似乎发生在 CoGBKResults 的覆盖方法中。

String Ad_ID = e.getKey();
String Ad_Info = "none";
Ad_Info = e.getValue().getOnly(AdInfoTag);

以下是我的加入方法。

static PCollection<String> joinEvents(PCollection<TableRow> ImpressionTable,
      PCollection<TableRow> AdTable) throws Exception {

    final TupleTag<String> ImpressionInfoTag = new TupleTag<String>();
    final TupleTag<String> AdInfoTag = new TupleTag<String>();

    // transform both input collections to tuple collections, where the keys are Ad_ID
    PCollection<KV<String, String>> ImpressionInfo = ImpressionTable.apply(
        ParDo.of(new ExtractImpressionDataInfoFn()));
    PCollection<KV<String, String>> AdInfo = AdTable.apply(
        ParDo.of(new ExtractAdDataInfoFn()));

    // Ad_ID 'key' -> CGBKR (<ImpressionInfo>, <AdInfo>)
    PCollection<KV<String, CoGbkResult>> kvpCollection = KeyedPCollectionTuple
        .of(ImpressionInfoTag, ImpressionInfo)
        .and(AdInfoTag, AdInfo)
        .apply(CoGroupByKey.<String>create());

    // Process the CoGbkResult elements generated by the CoGroupByKey transform.
    // Ad_ID 'key' -> string of <Impressioninfo>, <Adinfo>
    PCollection<KV<String, String>> finalResultCollection =
      kvpCollection.apply(ParDo.named("Process").of(
        new DoFn<KV<String, CoGbkResult>, KV<String, String>>() {
            private static final long serialVersionUID = 1L;

        @Override
          public void processElement(ProcessContext c) {
            KV<String, CoGbkResult> e = c.element();
            String Ad_ID = e.getKey();
            String Ad_Info = "none";
            Ad_Info = e.getValue().getOnly(AdInfoTag);
            for (String eventInfo : c.element().getValue().getAll(ImpressionInfoTag)) {
              // Generate a string that combines information from both collection values
              c.output(KV.of(Ad_ID, " " + Ad_Info
                      + " " + eventInfo));
            }
          }
      }));

     //write to GCS
    PCollection<String> formattedResults = finalResultCollection
        .apply(ParDo.named("Format").of(new DoFn<KV<String, String>, String>() {
          @Override
          public void processElement(ProcessContext c) {
            String outputstring = "AdUnitID: " + c.element().getKey()
                + ", " + c.element().getValue();
            c.output(outputstring);
          }
        }));
    return formattedResults;
  }

我的 ExtractImpressionDataInfoFn 类和 ExtractAdDatInfoFn 类如下。

static class ExtractImpressionDataInfoFn extends DoFn<TableRow, KV<String, String>> {
    private static final long serialVersionUID = 1L;

    @Override
    public void processElement(ProcessContext c) {
        TableRow row = c.element();
        String Ad_ID = (String) row.get("AdUnitID");
        String User_ID = (String) row.get("UserID");
        String Client_ID = (String) row.get("ClientID");
        String Impr_Time = (String) row.get("GfpActivityAdEventTIme");
        String ImprInfo = "UserID: " + User_ID + ", ClientID: " + Client_ID + ", GfpActivityAdEventTIme: " + Impr_Time;
        c.output(KV.of(Ad_ID, ImprInfo));
    }
}


static class ExtractAdDataInfoFn extends DoFn<TableRow, KV<String, String>> {
    private static final long serialVersionUID = 1L;

    @Override
    public void processElement(ProcessContext c) {
        TableRow row = c.element();
        String Ad_ID = (String) row.get("AdUnitID");
        String Content_ID = (String) row.get("ContentID");
        String Pub_ID = (String) row.get("Publisher");
        String Add_Info = "ContentID: " + Content_ID + ", Publisher: " + Pub_ID;
        c.output(KV.of(Ad_ID, Add_Info));
    }
}

展示和广告的架构如下 印象: 广告单元 ID 用户身份 客户 ID
GfpActivityAdEventTIme

广告: 广告单元 ID 客户编号 发布者

enter image description here

【问题讨论】:

标签: java google-bigquery google-cloud-platform google-cloud-dataflow gcp


【解决方案1】:

该错误表明当您调用 getOnly 时,CoGroupByKey 有多个结果。特别是这一行:

Ad_Info = e.getValue().getOnly(AdInfoTag);

如果您将其更改为getAll(AdInfoTag),它应该可以工作。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2012-09-10
    • 1970-01-01
    • 1970-01-01
    • 2019-12-23
    • 1970-01-01
    • 1970-01-01
    • 2020-11-21
    相关资源
    最近更新 更多