【问题标题】:Read data from database using UDF in pig在 pig 中使用 UDF 从数据库中读取数据
【发布时间】:2017-09-05 01:31:09
【问题描述】:

我需要从数据库中读取数据并使用 pig 分析数据。 我用java写了一个UDF Referring following link

register /tmp/UDFJars/CassandraUDF_1-0.0.1-SNAPSHOT-jar-with-dependencies.jar;
A = Load '/user/sampleFile.txt' using udf.DBLoader('10.xx.xxx.4','username','password','select * from customer limit 10') as (f1 : chararray);
DUMP A;


package udf;

import java.io.IOException;
import java.util.ArrayList;
import java.util.List;

import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.apache.hadoop.mapreduce.InputFormat;
import org.apache.hadoop.mapreduce.Job;
import org.apache.hadoop.mapreduce.RecordReader;
import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
import org.apache.hadoop.mapreduce.lib.input.TextInputFormat;
import org.apache.pig.LoadFunc;
import org.apache.pig.backend.hadoop.executionengine.mapReduceLayer.PigSplit;
import org.apache.pig.data.Tuple;
import org.apache.pig.data.TupleFactory;

import com.data.ConnectionCassandra;
import com.datastax.driver.core.ResultSet;
import com.datastax.driver.core.Row;
import com.datastax.driver.core.Session;

public class DBLoader extends LoadFunc {
    private final Log log = LogFactory.getLog(getClass());
    Session session;
    private ArrayList mProtoTuple = null;
    private String jdbcURL;
    private String user;
    private String pass;
    private int count = 0;
    private String query;
    ResultSet result;
    List<Row> rows;
    int colSize;
    protected TupleFactory mTupleFactory = TupleFactory.getInstance();

    public DBLoader() {
    }

    public DBLoader(String jdbcURL, String user, String pass, String query) {

        this.jdbcURL = jdbcURL;
        this.user = user;
        this.pass = pass;
        this.query = query;

    }

    @Override
    public InputFormat getInputFormat() throws IOException {
        log.info("Inside InputFormat");
        // TODO Auto-generated method stub
        try {
            return new TextInputFormat();
        } catch (Exception exception) {
            log.error(exception.getMessage());
            log.error(exception.fillInStackTrace());
            throw new IOException();
        }
    }

    @Override
    public Tuple getNext() throws IOException {
        log.info("Inside get Next");
        Row row = rows.get(count);
        if (row != null) {
            mProtoTuple = new ArrayList<Object>();
            for (int colNum = 0; colNum < colSize; colNum++) {
                mProtoTuple.add(row.getObject(colNum));
            }
        } else {
            return null;
        }
        Tuple t = mTupleFactory.newTuple(mProtoTuple);
        mProtoTuple.clear();
        return t;

    }

    @Override
    public void prepareToRead(RecordReader arg0, PigSplit arg1) throws IOException {
        log.info("Inside Prepare to Read");
        session = null;
        if (query == null) {
            throw new IOException("SQL Insert command not specified");
        }
        if (user == null || pass == null) {
            log.info("Creating Session with user name and password as: " + user + " : " + pass);
            session = ConnectionCassandra.connectToCassandra1(jdbcURL, user, pass);
            log.info("Session Created");
        } else {
            session = ConnectionCassandra.connectToCassandra1(jdbcURL, user, pass);
        }
        log.info("Executing Query " + query);
        result = session.execute(query);
        log.info("Query Executed :" + query);
        rows = result.all();
        count = 0;
        colSize = result.getColumnDefinitions().asList().size();
    }

    @Override
    public void setLocation(String location, Job job) throws IOException {
        log.info("Inside Set Location");
        try {
            FileInputFormat.setInputPaths(job, location);
        } catch (Exception exception) {
            log.info("Some thing went wrong : " + exception.getMessage());
            log.debug(exception);
        }

    }
}

上面是我的猪脚本和java代码。 这里的 /user/sampleFile.txt 是一个没有数据的虚拟文件。 我收到以下异常:

猪栈跟踪

错误 1066:无法打开别名 A 的迭代器

org.apache.pig.impl.logicalLayer.FrontendException:错误 1066:无法打开别名 A 的迭代器 在 org.apache.pig.PigServer.openIterator(PigServer.java:892) 在 org.apache.pig.tools.grunt.GruntParser.processDump(GruntParser.java:774) 在 org.apache.pig.tools.pigscript.parser.PigScriptParser.parse(PigScriptParser.java:372) 在 org.apache.pig.tools.grunt.GruntParser.parseStopOnError(GruntParser.java:198) 在 org.apache.pig.tools.grunt.GruntParser.parseStopOnError(GruntParser.java:173) 在 org.apache.pig.tools.grunt.Grunt.exec(Grunt.java:84) 在 org.apache.pig.Main.run(Main.java:484) 在 org.apache.pig.Main.main(Main.java:158) 在 sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) 在 sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:57) 在 sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) 在 java.lang.reflect.Method.invoke(Method.java:606) 在 org.apache.hadoop.util.RunJar.run(RunJar.java:221) 在 org.apache.hadoop.util.RunJar.main(RunJar.java:136) 原因:java.io.IOException:作业以异常状态FAILED终止 在 org.apache.pig.PigServer.openIterator(PigServer.java:884)

... 13 更多

【问题讨论】:

  • 猪码在哪里?从错误来看,您的 pig 陈述不正确。
  • @ANI register /tmp/UDFJars/CassandraUDF_1-0.0.1-SNAPSHOT-jar-with-dependencies.jar; A = Load '/user/sampleFile.txt' using udf.DBLoader('10.xx.xxx.4','username','password','select * from customer limit 10') as (f1 : chararray);转储 A;是我的简单猪脚本。
  • select * from customer limit 10 存储在一个变量 f1:chararray 中!表格是否只有一列?
  • 否,但我也尝试更改查询以从客户限制 10 中选择 customer_id 仍然得到相同的错误,我也很想知道这种获取数据的方法是否正确。
  • 另一种方法是使用 Sqoop 将数据读入 Hive,然后使用 HCatalgo 将数据读入 Pig。这些工具已经存在,没有必要重新发明轮子。如果你必须创建一个 UDF,那么只需在你的 UDF 中嵌入 Sqoop 导入命令。

标签: apache-pig udf


【解决方案1】:

维韦克!你甚至进入 prepareToRead 吗? (我看到你做了一些日志记录,所以很高兴知道你在日志中实际有什么)另外,提供完整的堆栈跟踪真的很棒,因为我看到你没有完整的底层异常。 只是一些想法 - 我从来没有尝试在没有实现自己的 InputFormat 和 RecordReader 的情况下编写 LoadFunc - TextInputFormat 检查文件是否存在及其大小(并根据文件大小创建多个 InputSplits),所以如果你的虚拟文件在那里是空的很可能没有产生 InputSplit 或产生零长度 InputSplit。由于它的长度为零,因此可能会导致 pig 抛出该异常。所以好的建议是实现自己的 InputFormat(实际上很容易)。也可以作为快速尝试 - 尝试

set pig.splitCombination false

可能没有帮助,但很容易尝试。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2020-02-14
    • 1970-01-01
    • 1970-01-01
    • 2018-02-17
    • 1970-01-01
    相关资源
    最近更新 更多