Class TextTable

  • All Implemented Interfaces:
    java.io.Serializable, BeamSqlTable
    Direct Known Subclasses:
    TextJsonTable

    @Internal
    public class TextTable
    extends SchemaBaseBeamTable
    TextTable is a BeamSqlTable that reads text files and converts them according to the specified format.

    Support formats are "csv" and "lines".

    CSVFormat itself has many dialects, check its javadoc for more info.

    See Also:
    Serialized Form
    • Constructor Summary

      Constructors 
      Constructor Description
      TextTable​(org.apache.beam.sdk.schemas.Schema schema, java.lang.String filePattern, org.apache.beam.sdk.transforms.PTransform<org.apache.beam.sdk.values.PCollection<java.lang.String>,​org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.values.Row>> readConverter, org.apache.beam.sdk.transforms.PTransform<org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.values.Row>,​org.apache.beam.sdk.values.PCollection<java.lang.String>> writeConverter)
      Text table with the specified read and write transforms.
    • Method Summary

      All Methods Instance Methods Concrete Methods 
      Modifier and Type Method Description
      org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.values.Row> buildIOReader​(org.apache.beam.sdk.values.PBegin begin)
      create a PCollection<Row> from source.
      org.apache.beam.sdk.values.PDone buildIOWriter​(org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.values.Row> input)
      create a IO.write() instance to write to target.
      java.lang.String getFilePattern()  
      BeamTableStatistics getTableStatistics​(org.apache.beam.sdk.options.PipelineOptions options)
      Estimates the number of rows or the rate for unbounded Tables.
      org.apache.beam.sdk.values.PCollection.IsBounded isBounded()
      Whether this table is bounded (known to be finite) or unbounded (may or may not be finite).
      • Methods inherited from class java.lang.Object

        clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
    • Constructor Detail

      • TextTable

        public TextTable​(org.apache.beam.sdk.schemas.Schema schema,
                         java.lang.String filePattern,
                         org.apache.beam.sdk.transforms.PTransform<org.apache.beam.sdk.values.PCollection<java.lang.String>,​org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.values.Row>> readConverter,
                         org.apache.beam.sdk.transforms.PTransform<org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.values.Row>,​org.apache.beam.sdk.values.PCollection<java.lang.String>> writeConverter)
        Text table with the specified read and write transforms.
    • Method Detail

      • getFilePattern

        public java.lang.String getFilePattern()
      • getTableStatistics

        public BeamTableStatistics getTableStatistics​(org.apache.beam.sdk.options.PipelineOptions options)
        Description copied from interface: BeamSqlTable
        Estimates the number of rows or the rate for unbounded Tables. If it is not possible to estimate the row count or rate it will return BeamTableStatistics.BOUNDED_UNKNOWN.
        Specified by:
        getTableStatistics in interface BeamSqlTable
        Overrides:
        getTableStatistics in class BaseBeamTable
      • isBounded

        public org.apache.beam.sdk.values.PCollection.IsBounded isBounded()
        Description copied from interface: BeamSqlTable
        Whether this table is bounded (known to be finite) or unbounded (may or may not be finite).
      • buildIOReader

        public org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.values.Row> buildIOReader​(org.apache.beam.sdk.values.PBegin begin)
        Description copied from interface: BeamSqlTable
        create a PCollection<Row> from source.
      • buildIOWriter

        public org.apache.beam.sdk.values.PDone buildIOWriter​(org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.values.Row> input)
        Description copied from interface: BeamSqlTable
        create a IO.write() instance to write to target.