Class ParquetDataWriter
java.lang.Object
com.northconcepts.datapipeline.core.DataObject
com.northconcepts.datapipeline.core.Endpoint
com.northconcepts.datapipeline.core.DataEndpoint
com.northconcepts.datapipeline.core.DataWriter
com.northconcepts.datapipeline.internal.lang.IntegrationWriter
com.northconcepts.datapipeline.parquet.ParquetDataWriter
Writes records to Apache Parquet columnar files.
See Apache Parquet columnar storage.
-
Nested Class Summary
Nested classes/interfaces inherited from class com.northconcepts.datapipeline.core.DataEndpoint
DataEndpoint.State -
Field Summary
Fields inherited from class com.northconcepts.datapipeline.core.DataEndpoint
lastRecord, PRODUCT, PRODUCT_VERSION, VENDOR, XML_INPUT_FACTORY_KEYFields inherited from class com.northconcepts.datapipeline.core.Endpoint
BUFFER_SIZE, captureElapsedTime, DEFAULT_READ_BUFFER_SIZEFields inherited from class com.northconcepts.datapipeline.core.DataObject
id, log, name, TIMESTAMP_FORMAT -
Constructor Summary
ConstructorsConstructorDescriptionParquetDataWriter(File file) Write parquet data to a file.ParquetDataWriter(OutputFile outputFile) Write parquet data to anOutputFile. -
Method Summary
Modifier and TypeMethodDescriptionaddExceptionProperties(DataException exception) Adds this endpoint's current state to aDataException.voidclose()Indicates that this endpoint has finished reading or writing.Indicates the folder to store cached files during dynamic schema generation withLocalFileDatasetor null if schema generation should be performed completely in memory usingMemoryDataset/> See alsogetRecordsPerCacheFile()protected intReturns the number of threads that collect column statistics when the schema is inferred from the records (default is 2).CompressionCodecNameIndicates the compression used for writing (default UNCOMPRESSED).ConfigurationReturns the Parquet configuration parameters.intReturns the default scale used when writing BigDecimal values (default 5).intReturns the default precision used when writing BigDecimal invalid input: '&' BigInteger values (default 25).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).intIndicates 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).Indicates the rounding algorithm used for all BigDecimal values (default isRoundingMode.HALF_UP).MessageTypeReturns the schema used to write the file.booleanDeprecated.booleanIndicates if all datetime fields should be marked as AdjustedToUTC.booleanIndicates if unsupported characters should be removed from field names (default is true).voidopen()Makes this endpoint ready for reading or writing.setCacheFolder(String cacheFolder) Indicates the folder to store cached files during dynamic schema generation withLocalFileDatasetor null if schema generation should be performed completely in memory usingMemoryDataset.
See alsosetRecordsPerCacheFile(int)protected ParquetDataWritersetColumnStatsReaderThreads(int columnStatsReaderThreads) Sets the number of threads that collect column statistics when the schema is inferred from the records (default is 2).setCompressionCodecName(CompressionCodecName compressionCodecName) Indicates the compression used for writing (default UNCOMPRESSED).setConfiguration(Configuration configuration) Sets the Parquet configuration parameters.setDefaultAdjustedToUTC(boolean defaultAdjustedToUTC) Indicates if all datetime fields should be marked as AdjustedToUTC.setDefaultAdjustToUTC(boolean defaultAdjustToUTC) Deprecated.setDefaultBigDecimalScale(int defaultBigDecimalScale) Sets the default scale used when writing BigDecimal values (default 5).setDefaultBigNumberPrecision(int defaultBigNumberPrecision) Sets the default precision used when writing BigDecimal invalid input: '&' BigInteger values (default 25).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).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).setRemoveUnsupportedChars(boolean removeUnsupportedChars) Indicates if unsupported characters should be removed from field names (default is true).setRoundingMode(RoundingMode roundingMode) Indicates the rounding algorithm used for all BigDecimal values (default isRoundingMode.HALF_UP).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(JdbcConnectionFactory jdbcConnectionFactory, String query, Object... queryParameters) Sets the schema used to write the file by copying it from the metadata of an SQL query.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.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.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(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(MessageType schema) Sets the schema used to write the file.protected voidOverridden by subclasses to write the specified record to thisDataWriter.Methods inherited from class com.northconcepts.datapipeline.core.DataWriter
available, getNestedEndpoint, getNestedWriter, getRootEndpoint, getRootWriter, getWriter, writeMethods inherited from class com.northconcepts.datapipeline.core.DataEndpoint
decrementRecordCount, enableJmx, getLastRecord, getRecordCount, getRecordCountAsBigInteger, getRecordCountAsString, incrementRecordCount, isRecordCountBigInteger, resetRecordCount, toStringMethods inherited from class com.northconcepts.datapipeline.core.Endpoint
addElapsedtime, assertClosed, assertNotOpened, assertOpened, finalize, getClosedOn, getDescription, getElapsedTime, getElapsedTimeAsString, getOpenedOn, getOpenElapsedTime, getOpenElapsedTimeAsString, getSelfTime, getSelfTimeAsString, getState, isCaptureElapsedTime, isClosed, isOpen, setCaptureElapsedTime, setDescription
-
Constructor Details
-
ParquetDataWriter
Write parquet data to a file. -
ParquetDataWriter
public ParquetDataWriter(OutputFile outputFile) Write parquet data to anOutputFile.- Parameters:
outputFile- - OutputFile with FileSystem.
-
-
Method Details
-
open
Description copied from class:DataEndpointMakes this endpoint ready for reading or writing.- Overrides:
openin classIntegrationWriter- Throws:
DataException
-
writeImpl
Description copied from class:DataWriterOverridden by subclasses to write the specified record to thisDataWriter.Contract for subclasses (see also
docs/authoring/DataWriter.md):- Do not call this method directly —
DataWriter.write(Record)is the template method that wraps exceptions, increments the record count, and attaches the offending record to thrownDataExceptions. - Do not call
DataEndpoint.incrementRecordCount()here;DataWriter.write(Record)already does. - Throw raw exceptions;
DataWriter.write(Record)wraps them viaexception(throwable).setRecord(record).
- Specified by:
writeImplin classDataWriter- Throws:
Throwable
- Do not call this method directly —
-
close
Description copied from class:DataEndpointIndicates that this endpoint has finished reading or writing.- Overrides:
closein classDataEndpoint- Throws:
DataException
-
getSchema
public MessageType getSchema()Returns the schema used to write the file. -
setSchema
Sets the schema used to write the file. -
setSchema
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
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
Sets the default precision used when writing BigDecimal invalid input: '&' BigInteger values (default 25). -
addExceptionProperties
Description copied from class:EndpointAdds this endpoint's current state to aDataException. Since this method is called whenever an exception is thrown, subclasses should override it to add their specific information.- Overrides:
addExceptionPropertiesin classDataWriter
-
getCompressionCodecName
public CompressionCodecName getCompressionCodecName()Indicates the compression used for writing (default UNCOMPRESSED). -
setCompressionCodecName
Indicates the compression used for writing (default UNCOMPRESSED). -
isDefaulAdjustToUTC
Deprecated. -
setDefaultAdjustToUTC
Deprecated. -
isDefaultAdjustedToUTC
public boolean isDefaultAdjustedToUTC()Indicates if all datetime fields should be marked as AdjustedToUTC. -
setDefaultAdjustedToUTC
Indicates if all datetime fields should be marked as AdjustedToUTC. -
getRoundingMode
Indicates the rounding algorithm used for all BigDecimal values (default isRoundingMode.HALF_UP). -
setRoundingMode
Indicates the rounding algorithm used for all BigDecimal values (default isRoundingMode.HALF_UP). -
getConfiguration
public Configuration getConfiguration()Returns the Parquet configuration parameters. -
setConfiguration
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
Sets the number of threads that collect column statistics when the schema is inferred from the records (default is 2). -
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 innullwill 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: Usingnullor a high record count can significantly slow down processing and cause anOutOfMemoryError. -
setMaxRecordsAnalyzed
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 innullwill 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: Usingnullor a high record count can significantly slow down processing and cause anOutOfMemoryError. -
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 ifsetCacheFolder(String)is set. -
setRecordsPerCacheFile
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 ifsetCacheFolder(String)is set. -
getCacheFolder
Indicates the folder to store cached files during dynamic schema generation withLocalFileDatasetor null if schema generation should be performed completely in memory usingMemoryDataset/> See alsogetRecordsPerCacheFile() -
setCacheFolder
Indicates the folder to store cached files during dynamic schema generation withLocalFileDatasetor null if schema generation should be performed completely in memory usingMemoryDataset.
See alsosetRecordsPerCacheFile(int) -
isRemoveUnsupportedChars
public boolean isRemoveUnsupportedChars()Indicates if unsupported characters should be removed from field names (default is true). -
setRemoveUnsupportedChars
Indicates if unsupported characters should be removed from field names (default is true).
-
isDefaultAdjustedToUTC()