Skip to content

Instantly share code, notes, and snippets.

@damondouglas
Created September 25, 2023 17:25
Show Gist options
  • Select an option

  • Save damondouglas/84c34469d4a6d7b7483468659482855c to your computer and use it in GitHub Desktop.

Select an option

Save damondouglas/84c34469d4a6d7b7483468659482855c to your computer and use it in GitHub Desktop.
Example Skipping Lines using Apache Beam FileIO ReadableFile
/*
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance
* with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.apache.beam.examples;
import static org.apache.beam.sdk.util.Preconditions.checkStateNotNull;
import avro.shaded.com.google.common.base.Throwables;
import java.io.IOException;
import java.util.List;
import java.util.stream.Collectors;
import java.util.stream.Stream;
import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.io.FileIO;
import org.apache.beam.sdk.io.FileIO.ReadableFile;
import org.apache.beam.sdk.io.fs.ResourceId;
import org.apache.beam.sdk.transforms.DoFn;
import org.apache.beam.sdk.transforms.PTransform;
import org.apache.beam.sdk.transforms.ParDo;
import org.apache.beam.sdk.values.KV;
import org.apache.beam.sdk.values.PCollection;
import org.apache.beam.sdk.values.PCollectionTuple;
import org.apache.beam.sdk.values.TupleTag;
import org.apache.beam.sdk.values.TupleTagList;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
public class FileIOExample {
private static final Logger LOG = LoggerFactory.getLogger(FileIOExample.class);
private static final TupleTag<KV<String, List<String>>> OUTPUT =
new TupleTag<KV<String, List<String>>>() {};
private static final TupleTag<KV<String, String>> FAILURE = new TupleTag<KV<String, String>>() {};
public static void main(String[] args) {
Pipeline p = Pipeline.create();
PCollectionTuple readLinesResult =
p.apply("match", FileIO.match().filepattern("gs://apache-beam-samples/shakespeare/*.txt"))
.apply("readMatches", FileIO.readMatches())
.apply("readLines", new ReadLines(100));
readLinesResult
.get(FAILURE)
.apply(
"logFailures",
ParDo.of(
new DoFn<KV<String, String>, Void>() {
@ProcessElement
public void process(@Element KV<String, String> element) {
LOG.error("{}: {}", element.getKey(), element.getValue());
}
}));
readLinesResult
.get(OUTPUT)
.apply("head", ParDo.of(new HeadFn()))
.apply(
"logOutput",
ParDo.of(
new DoFn<String, Void>() {
@ProcessElement
public void process(@Element String head) {
LOG.info(head);
}
}));
p.run();
}
static class ReadLines extends PTransform<PCollection<ReadableFile>, PCollectionTuple> {
private final int skipLines;
ReadLines(int skipLines) {
this.skipLines = skipLines;
}
@Override
public PCollectionTuple expand(PCollection<ReadableFile> input) {
return input.apply(
ReadLinesFn.class.getSimpleName(),
ParDo.of(new ReadLinesFn(this)).withOutputTags(OUTPUT, TupleTagList.of(FAILURE)));
}
}
static class ReadLinesFn extends DoFn<ReadableFile, KV<String, List<String>>> {
private final ReadLines spec;
ReadLinesFn(ReadLines spec) {
this.spec = spec;
}
@ProcessElement
public void process(@Element ReadableFile element, MultiOutputReceiver receiver) {
ResourceId resourceId = element.getMetadata().resourceId();
try {
String rawData = element.readFullyAsUTF8String();
Stream<String> lines = rawData.lines().skip(spec.skipLines);
List<String> result = lines.collect(Collectors.toList());
receiver.get(OUTPUT).output(KV.of(resourceId.toString(), result));
} catch (IOException e) {
receiver
.get(FAILURE)
.output(KV.of(resourceId.toString(), Throwables.getStackTraceAsString(e)));
}
}
}
static class HeadFn extends DoFn<KV<String, List<String>>, String> {
@ProcessElement
public void process(
@Element KV<String, List<String>> element, OutputReceiver<String> receiver) {
String key = checkStateNotNull(element.getKey());
List<String> lines = checkStateNotNull(element.getValue());
int to = 5;
if (lines.size() < to) {
to = lines.size();
}
for (int i = 0; i < to; i++) {
String line = lines.get(i);
receiver.output(String.format("%s[%d]: %s", key, i + 1, line));
}
}
}
}
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment