View difference between Paste ID: XQjAHNY5 and Vpa88Qbp
SHOW: | | - or go back to the newest paste.
1
public class MyLoadFunc extends LoadFunc {
2
    protected RecordReader in = null;
3
    private byte fieldDel = '\t';
4
    private ArrayList<Object> mProtoTuple = null;
5
    private TupleFactory mTupleFactory = TupleFactory.getInstance();
6
    private static final int BUFFER_SIZE = 1024;
7
    private static Log log = LogFactory.getLog(GeofenceLoadFunc.class);
8
9
    @Override
10
    public Tuple getNext() throws IOException {
11
        try {
12
            boolean notDone = in.nextKeyValue();
13
14
            if (!notDone) {
15
                return null;
16
            }
17
            Text value = (Text) in.getCurrentValue();
18
19
            byte[] buf = value.getBytes();
20
            int len = value.getLength();
21
            int start = 0;
22
23
            for (int i = 0; i < len; i++) {
24
                if (buf[i] == fieldDel) {
25
                    readField(buf, start, i);
26
                    start = i + 1;
27
                }
28
            }
29
            // pick up the last field
30
            readField(buf, start, len);
31
32
            Tuple t =  mTupleFactory.newTupleNoCopy(mProtoTuple);
33
            mProtoTuple = null;
34
            return t;
35
        } catch (InterruptedException e) {
36
            int errCode = 6018;
37
            String errMsg = "Error while reading input";
38
            throw new ExecException(errMsg, errCode,
39
                    PigException.REMOTE_ENVIRONMENT, e);
40
        }
41
42
    }
43
44
    private void readField(byte[] buf, int start, int end) {
45
        if (mProtoTuple == null) {
46
            mProtoTuple = new ArrayList<Object>();
47
        }
48
49
        if (start == end) {
50
            // NULL value
51
            mProtoTuple.add(null);
52
        } else {
53
            mProtoTuple.add(new DataByteArray(buf, start, end));
54
        }
55
    }
56
57
    @Override
58
    public InputFormat getInputFormat() {
59
        return new TextInputFormat();
60
    }
61
62
    @Override
63
    public void prepareToRead(RecordReader reader, PigSplit split) {
64
        in = reader;
65
    }
66
67
    @Override
68
    public void setLocation(String location, Job job)
69
            throws IOException {
70
        FileInputFormat.setInputPaths(job, location);
71
    }
72
}