Class AvroSinks

java.lang.Object
com.hazelcast.jet.avro.AvroSinks

public final class AvroSinks extends Object
Contains factory methods for Apache Avro sinks.
Since:
Jet 3.0
  • Method Summary

    Modifier and Type
    Method
    Description
    static <R> com.hazelcast.jet.pipeline.Sink<R>
    files(String directoryName, Class<R> recordClass, org.apache.avro.Schema schema)
    Convenience for files(String, Schema, SupplierEx) which uses either SpecificDatumWriter or ReflectDatumWriter depending on the supplied recordClass.
    static com.hazelcast.jet.pipeline.Sink<org.apache.avro.generic.IndexedRecord>
    files(String directoryName, org.apache.avro.Schema schema)
    Convenience for files(String, Schema, SupplierEx) which uses GenericDatumWriter.
    static <R> com.hazelcast.jet.pipeline.Sink<R>
    files(String directoryName, org.apache.avro.Schema schema, com.hazelcast.function.SupplierEx<org.apache.avro.io.DatumWriter<R>> datumWriterSupplier)
    Returns a sink that that writes the items it receives to Apache Avro files.

    Methods inherited from class java.lang.Object

    clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
  • Method Details

    • files

      @Nonnull public static <R> com.hazelcast.jet.pipeline.Sink<R> files(@Nonnull String directoryName, @Nonnull org.apache.avro.Schema schema, @Nonnull com.hazelcast.function.SupplierEx<org.apache.avro.io.DatumWriter<R>> datumWriterSupplier)
      Returns a sink that that writes the items it receives to Apache Avro files. Each processor will write to its own file whose name is equal to the processor's global index (an integer unique to each processor of the vertex), but a single pathname is used to resolve the containing directory of all files, on all cluster members. The sink always overwrites the files.

      The sink creates a DataFileWriter for each processor using the supplied datumWriterSupplier with the given Schema.

      No state is saved to snapshot for this sink. After the job is restarted, the items will be missing since files will be overwritten.

      The default local parallelism for this sink is 1.

      Type Parameters:
      R - the type of the record
      Parameters:
      directoryName - directory to create the files in. Will be created if it doesn't exist. Must be the same on all members.
      schema - the record schema
      datumWriterSupplier - the record writer supplier
    • files

      @Nonnull public static <R> com.hazelcast.jet.pipeline.Sink<R> files(@Nonnull String directoryName, @Nonnull Class<R> recordClass, @Nonnull org.apache.avro.Schema schema)
      Convenience for files(String, Schema, SupplierEx) which uses either SpecificDatumWriter or ReflectDatumWriter depending on the supplied recordClass.
    • files

      @Nonnull public static com.hazelcast.jet.pipeline.Sink<org.apache.avro.generic.IndexedRecord> files(@Nonnull String directoryName, @Nonnull org.apache.avro.Schema schema)
      Convenience for files(String, Schema, SupplierEx) which uses GenericDatumWriter.