Class ParquetDataWriter


public class ParquetDataWriter extends IntegrationWriter
Writes records to Apache Parquet columnar files. See Apache Parquet columnar storage.
  • Constructor Details

    • ParquetDataWriter

      public ParquetDataWriter(File file)
      Write parquet data to a file.
    • ParquetDataWriter

      public ParquetDataWriter(OutputFile outputFile)
      Write parquet data to an OutputFile.
      Parameters:
      outputFile - - OutputFile with FileSystem.
  • Method Details

    • open

      public void open() throws DataException
      Description copied from class: DataEndpoint
      Makes this endpoint ready for reading or writing.
      Overrides:
      open in class IntegrationWriter
      Throws:
      DataException
    • writeImpl

      protected void writeImpl(Record record) throws Throwable
      Description copied from class: DataWriter
      Overridden by subclasses to write the specified record to this DataWriter.

      Contract for subclasses (see also docs/authoring/DataWriter.md):

      Specified by:
      writeImpl in class DataWriter
      Throws:
      Throwable
    • close

      public void close() throws DataException
      Description copied from class: DataEndpoint
      Indicates that this endpoint has finished reading or writing.
      Overrides:
      close in class DataEndpoint
      Throws:
      DataException
    • getSchema

      public MessageType getSchema()
      Returns the schema used to write the file.
    • setSchema

      public ParquetDataWriter setSchema(MessageType schema)
      Sets the schema used to write the file.
    • setSchema

      public ParquetDataWriter setSchema(Connection connection, String query, Object... queryParameters)
      Sets the schema used to write the file by copying it from the metadata of an SQL query.
    • setSchema

      public ParquetDataWriter setSchema(Connection connection, JdbcValueReader jdbcValueReader, String query, Object... queryParameters)
      Sets the schema used to write the file by copying it from the metadata of an SQL query. Since no data is read from query (only its metadata is used), care should be taken to create an optimized query where the database does minimal work and returns no data. For example, consider using a query like: SELECT * FROM invoices WHERE 1invalid input: '<'0.
    • setSchema

      public ParquetDataWriter setSchema(JdbcConnectionFactory jdbcConnectionFactory, String query, Object... queryParameters)
      Sets the schema used to write the file by copying it from the metadata of an SQL query. Since no data is read from query (only its metadata is used), care should be taken to create an optimized query where the database does minimal work and returns no data. For example, consider using a query like: SELECT * FROM invoices WHERE 1invalid input: '<'0.
    • setSchema

      public ParquetDataWriter setSchema(JdbcConnectionFactory jdbcConnectionFactory, JdbcValueReader jdbcValueReader, String query, Object... queryParameters)
      Sets the schema used to write the file by copying it from the metadata of an SQL query.
    • setSchema

      public ParquetDataWriter setSchema(Connection connection, String databaseCatalog, String databaseSchema, String databaseTable)
      Sets the schema used to write the Parquet file by copying it from the schema of a database table.
    • setSchema

      public ParquetDataWriter setSchema(Connection connection, JdbcValueReader sqlToJavaTypeMapper, String databaseCatalog, String databaseSchema, String databaseTable)
      Sets the schema used to write the Parquet file by copying it from the schema of a database table.
    • getDefaultBigDecimalScale

      public int getDefaultBigDecimalScale()
      Returns the default scale used when writing BigDecimal values (default 5).
    • setDefaultBigDecimalScale

      public ParquetDataWriter setDefaultBigDecimalScale(int defaultBigDecimalScale)
      Sets the default scale used when writing BigDecimal values (default 5).
    • getDefaultBigNumberPrecision

      public int getDefaultBigNumberPrecision()
      Returns the default precision used when writing BigDecimal invalid input: '&' BigInteger values (default 25).
    • setDefaultBigNumberPrecision

      public ParquetDataWriter setDefaultBigNumberPrecision(int defaultBigNumberPrecision)
      Sets the default precision used when writing BigDecimal invalid input: '&' BigInteger values (default 25).
    • addExceptionProperties

      public DataException addExceptionProperties(DataException exception)
      Description copied from class: Endpoint
      Adds this endpoint's current state to a DataException. Since this method is called whenever an exception is thrown, subclasses should override it to add their specific information.
      Overrides:
      addExceptionProperties in class DataWriter
    • getCompressionCodecName

      public CompressionCodecName getCompressionCodecName()
      Indicates the compression used for writing (default UNCOMPRESSED).
    • setCompressionCodecName

      public ParquetDataWriter setCompressionCodecName(CompressionCodecName compressionCodecName)
      Indicates the compression used for writing (default UNCOMPRESSED).
    • isDefaulAdjustToUTC

      @Deprecated public boolean isDefaulAdjustToUTC()
      Deprecated.
    • setDefaultAdjustToUTC

      @Deprecated public ParquetDataWriter setDefaultAdjustToUTC(boolean defaultAdjustToUTC)
    • isDefaultAdjustedToUTC

      public boolean isDefaultAdjustedToUTC()
      Indicates if all datetime fields should be marked as AdjustedToUTC.
    • setDefaultAdjustedToUTC

      public ParquetDataWriter setDefaultAdjustedToUTC(boolean defaultAdjustedToUTC)
      Indicates if all datetime fields should be marked as AdjustedToUTC.
    • getRoundingMode

      public RoundingMode getRoundingMode()
      Indicates the rounding algorithm used for all BigDecimal values (default is RoundingMode.HALF_UP).
    • setRoundingMode

      public ParquetDataWriter setRoundingMode(RoundingMode roundingMode)
      Indicates the rounding algorithm used for all BigDecimal values (default is RoundingMode.HALF_UP).
    • getConfiguration

      public Configuration getConfiguration()
      Returns the Parquet configuration parameters.
    • setConfiguration

      public ParquetDataWriter setConfiguration(Configuration configuration)
      Sets the Parquet configuration parameters.
    • getColumnStatsReaderThreads

      protected int getColumnStatsReaderThreads()
      Returns the number of threads that collect column statistics when the schema is inferred from the records (default is 2).
    • setColumnStatsReaderThreads

      protected ParquetDataWriter setColumnStatsReaderThreads(int columnStatsReaderThreads)
      Sets the number of threads that collect column statistics when the schema is inferred from the records (default is 2).
    • getMaxRecordsAnalyzed

      public Long getMaxRecordsAnalyzed()
      Indicates how many records should be analyzed and cached to generate the Parquet schema if no schema was explicitly set on this writer (default is 1000). This value will not be used if a schema was set on this writer.

      Passing in null will cause all records to be read and cached to determine the schema.

      The value will be set to 1 if a value less than 1 is passed in.

      Note: Using null or a high record count can significantly slow down processing and cause an OutOfMemoryError.
    • setMaxRecordsAnalyzed

      public ParquetDataWriter setMaxRecordsAnalyzed(Long maxRecordsAnalyzed)
      Indicates how many records should be analyzed and cached to generate the Parquet schema if no schema was explicitly set on this writer (default is 1000). This value will not be used if a schema was set on this writer.

      Passing in null will cause all records to be read and cached to determine the schema.

      The value will be set to 1 if a value less than 1 is passed in.

      Note: Using null or a high record count can significantly slow down processing and cause an OutOfMemoryError.
    • getRecordsPerCacheFile

      public int getRecordsPerCacheFile()
      Indicates how many records should be cached in memory to generate the Parquet schema if no schema was explicitly set on this writer (default is 10_000L). This value will not be used if a schema was set on this writer.

      The value will be set to 10_000L if a value less than 1 is passed in.

      Note: Using a small record count can significantly slow down processing due to many IO operations. And using a very huge record count can reduce the performance and consume more memory as it will hold more records in memory.

      Used only if setCacheFolder(String) is set.
    • setRecordsPerCacheFile

      public ParquetDataWriter setRecordsPerCacheFile(int recordsPerCacheFile)
      Indicates how many records should be cached in memory to generate the Parquet schema if no schema was explicitly set on this writer (default is 10_000L). This value will not be used if a schema was set on this writer.

      The value will be set to 10_000 if a value less than 1 is passed in.

      Note: Using a small record count can significantly slow down processing due to many IO operations. And using a very huge record count can reduce the performance and consume more memory as it will hold more records in memory.

      Used only if setCacheFolder(String) is set.
    • getCacheFolder

      public String getCacheFolder()
      Indicates the folder to store cached files during dynamic schema generation with LocalFileDataset or null if schema generation should be performed completely in memory using MemoryDataset /> See also getRecordsPerCacheFile()
    • setCacheFolder

      public ParquetDataWriter setCacheFolder(String cacheFolder)
      Indicates the folder to store cached files during dynamic schema generation with LocalFileDataset or null if schema generation should be performed completely in memory using MemoryDataset.
      See also setRecordsPerCacheFile(int)
    • isRemoveUnsupportedChars

      public boolean isRemoveUnsupportedChars()
      Indicates if unsupported characters should be removed from field names (default is true).
    • setRemoveUnsupportedChars

      public ParquetDataWriter setRemoveUnsupportedChars(boolean removeUnsupportedChars)
      Indicates if unsupported characters should be removed from field names (default is true).