Not a member of Pastebin yet?
Sign Up,
it unlocks many cool features!
- public class MyLoadFunc extends LoadFunc {
- protected RecordReader in = null;
- private byte fieldDel = '\t';
- private ArrayList<Object> mProtoTuple = null;
- private TupleFactory mTupleFactory = TupleFactory.getInstance();
- private static final int BUFFER_SIZE = 1024;
- private static Log log = LogFactory.getLog(GeofenceLoadFunc.class);
- @Override
- public Tuple getNext() throws IOException {
- try {
- boolean notDone = in.nextKeyValue();
- if (!notDone) {
- return null;
- }
- Text value = (Text) in.getCurrentValue();
- byte[] buf = value.getBytes();
- int len = value.getLength();
- int start = 0;
- for (int i = 0; i < len; i++) {
- if (buf[i] == fieldDel) {
- readField(buf, start, i);
- start = i + 1;
- }
- }
- // pick up the last field
- readField(buf, start, len);
- Tuple t = mTupleFactory.newTupleNoCopy(mProtoTuple);
- mProtoTuple = null;
- return t;
- } catch (InterruptedException e) {
- int errCode = 6018;
- String errMsg = "Error while reading input";
- throw new ExecException(errMsg, errCode,
- PigException.REMOTE_ENVIRONMENT, e);
- }
- }
- private void readField(byte[] buf, int start, int end) {
- if (mProtoTuple == null) {
- mProtoTuple = new ArrayList<Object>();
- }
- if (start == end) {
- // NULL value
- mProtoTuple.add(null);
- } else {
- mProtoTuple.add(new DataByteArray(buf, start, end));
- }
- }
- @Override
- public InputFormat getInputFormat() {
- return new TextInputFormat();
- }
- @Override
- public void prepareToRead(RecordReader reader, PigSplit split) {
- in = reader;
- }
- @Override
- public void setLocation(String location, Job job)
- throws IOException {
- FileInputFormat.setInputPaths(job, location);
- }
- }
Add Comment
Please, Sign In to add comment