【问题标题】:Data cleansing with cassandra and pig使用 cassandra 和 pig 进行数据清理
【发布时间】:2012-11-16 05:48:50
【问题描述】:

我想在 Pig 中比较两组数据。两者都具有相同的唯一 ID,第二组数据中的名称随机更改。逻辑如下:

  • 加载 empl1 原始数据
  • 加载 empl2 原始数据
  • 选择“名称不相同”和“emplno 相等”的行

我已经完成了:

A1=  LOAD 'cassandra://employees_pig1/employees_cf' USING CassandraStorage() AS (key, columns: bag {T: tuple(name, value)});

B1=  LOAD 'cassandra://employees_pig2/employees_cf' USING CassandraStorage() AS (key, columns: bag {T: tuple(name, value)});

A2 = FOREACH A1 GENERATE key, FLATTEN(columns);

B2 = FOREACH B1 GENERATE key as key2, FLATTEN(columns);

呵呵,不能在论坛发图片。这是说明 A2,B2 的链接 https://picasaweb.google.com/lh/photo/SU3QgKsbA4nmq83cdnhiVdMTjNZETYmyPJy0liipFm0?feat=directlink

现在需要一些帮助,我的处理方法正确吗?

C1 = join A2 by key, B2 by key2;

D1= filter C1 by A2.key==B2.key2 -- cannot do a A2.first_name!=B2.first_name;

想要选择“名称不同”和“emplno 相等”的行,但不完全确定如何操作。请指教。

谢谢你

更新: - 而不是加入我做了一个cogroup C3= COGROUP A2 by key, B2 by key2;

https://picasaweb.google.com/lh/photo/_lkEqW4BvIgbnZSHKDCJGNMTjNZETYmyPJy0liipFm0?feat=directlink

接下来,我正在考虑做

D1= FOREACH C3 GENERATE group, A2.first_name as fn1, B2.first_name as fn2

组返回所需的结果(即empno),但'A2.first_name,B2.first_name'不正确。需要知道如何访问 A2 和 B2 包/元组中的数据。

然后我就可以通过 fn1==fn2 进行过滤。

【问题讨论】:

    标签: hadoop cassandra apache-pig datastax-enterprise


    【解决方案1】:

    通过执行JOIN(至少是内部连接,这是您在上面所做的),您已经注意确保来自ABemplnos 相等。那么你所要做的就是过滤names是否相同。

    C1 = join A2 by key, B2 by key;
    D1 = filter C1 by A2::name != B2::name;
    

    【讨论】:

    • 谢谢你,但我试过了,但它不起作用。
      错误:无效的字段投影。模式中不存在投影字段 [A2::first_name]。 需要一些指针如何访问列内的数据:<pre> &gt;describe A2 A2:{key:chararray,columns::name:chararray,colums ::value:chararray}
    • A2 中的字段被称为name 时,你为什么要做A2::first_namedescribe C1join 之后,您将看到这些字段的名称,以便您正确地形成过滤器。
    【解决方案2】:

    已解决:)

    步骤: - 下载pygmalionhttps://github.com/jeromatron/pygmalion/downloads

    快速测试:

    register '/usr/share/dse/pygmalion/pygmalion-1.0.0.jar';
    define FromCassandraBag org.pygmalion.udf.FromCassandraBag();
    define ToCassandraBag org.pygmalion.udf.ToCassandraBag();
    
    A1=  LOAD 'cassandra://employees_pig1/employees_cf' USING CassandraStorage() AS (key,
    columns: bag {T: tuple(name, value)});
    B1=  LOAD 'cassandra://employees_pig2/employees_cf' USING CassandraStorage() AS (key, 
    columns: bag {T: tuple(name, value)});
    
    A2 = foreach A1 generate key,
    flatten(org.pygmalion.udf.FromCassandraBag('first_name', columns))
    as (first_name: chararray);
    
    B2 = foreach B1 generate key,
    flatten(org.pygmalion.udf.FromCassandraBag('first_name', columns))
    as (first_name: chararray);
    
    C1 = join A2 by key, B2 by key;
    D1= filter C1 BY A2::first_name != B2::first_name;
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2021-07-01
      • 2017-10-15
      • 2018-01-13
      • 1970-01-01
      • 2020-07-10
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多