Class Dataset
java.lang.Object
com.northconcepts.datapipeline.foundations.core.Bean
com.northconcepts.datapipeline.foundations.core.FoundationObject
com.northconcepts.datapipeline.foundations.pipeline.dataset.Dataset
- All Implemented Interfaces:
DataExceptionContributor,JsonSerializable,RecordSerializable,XmlSerializable,Closeable,Serializable,AutoCloseable,Iterable<Record>
- Direct Known Subclasses:
JdbcDataset,LocalFileDataset,MemoryDataset,MvStoreDataset
The base class for caching records produced by a
Pipeline or DataMappingPipeline. This class handles asynchronous loading
and column metadata and statistics.- See Also:
-
Nested Class Summary
Nested ClassesModifier and TypeClassDescriptionstatic classReads one summary record per column, with its name, counts, inferred type, lengths and sample value. -
Field Summary
Fields inherited from class com.northconcepts.datapipeline.foundations.core.FoundationObject
internalId, internalName, log, TIMESTAMP_FORMATFields inherited from interface com.northconcepts.datapipeline.core.RecordSerializable
SERIALIZED_CLASS_NAME, TYPEFields inherited from interface com.northconcepts.datapipeline.core.XmlSerializable
XML_SERIALIZED_CLASS_NAME -
Constructor Summary
ConstructorsConstructorDescriptionDataset(AbstractPipeline pipeline) Creates a dataset that caches the given pipeline's records onceload()is called. -
Method Summary
Modifier and TypeMethodDescriptionprotected ColumnAdds the field to the stats of its column, creating the column if needed; runs on a column stats worker thread.protected voidCalled during the data loading process after all the column stats have been loaded.protected abstract voidCalled at the end of the data loading process after all the records and column stats have been loaded.protected voidCalled during the data loading process after all the records have been loaded.protected abstract voidCalled at the start of the data loading process, but before any records or column stats have been loaded.voidGracefully terminate the asynchronous data loading and column stats calculation, waiting up to 10 seconds in total.voidclose()Returns a reader that emits one summary record per column of this dataset (seeDataset.ColumnsDataReader).Reads all data cached in this Dataset.createDataReader(long offset, int count) Reads a subset of data cached in this Dataset.protected abstract DataWriterWrites records to this dataset's cache after clearing it.protected voidfinalize()voidPerforms the given action for each record cached in this dataset.fromRecord(Record source) Loads this instance's state from a record and returnsthis(for fluid API call chaining).abstract ColumngetColumn(int index) abstract ColumnReturns the column with the given name or null if there is none.abstract longReturns the last exception thrown while collecting column stats or null if there was none.intThe number of threads to use to process column stats (default 2).Returns the exception that failed the current or last load or null if there was none.getJob()Returns the job loading this dataset's records or null before the first load and aftercancelLoad().The number of records to use when calculating column stats ornullfor all records (defaultnull).Defines the first N columns to analyze for stats for column stats.The maximum records to load parameter (maxRecordsToLoad) passed to the last call toload(Integer)orload(Integer, JobCallback).protected abstract ColumngetOrCreateColumn(String name, int index) Returns the named column, creating it with the given 0-based field index if needed; called while collecting column stats.abstract RecordgetRecord(long index) Returns the loaded record at the given 0-based index.abstract longReturns the number of records loaded into this dataset so far.getRecordList(long offset, int count) Get a subset of the records cached in this dataset.booleanIndicates if unique values in the dataset should be collected (defaultfalse).booleanReturntrueif all the column stats have been loaded.booleanReturntrueif all the records have been loaded and all the column stats have been loaded.booleanReturntrueif the records or column stats are currently being loaded.booleanIndicates if big decimals and big integers should be analyzed to determine their precision and scale (defaulttrue).booleanIndicates if boolean values should be looked for in strings and undefined types (defaulttrue).booleanIndicates if numeric values should be looked for in strings and undefined types (defaulttrue).booleanIndicates if date/time patterns should be looked for in strings and undefined types (defaulttrue).booleanIndicates if UUID values should be looked for in strings and undefined types (defaulttrue).booleanIndicates if string and undefined types should be analyzed to determine if they represent a numeric, boolean, or date/time value (defaulttrue).booleanReturntrueif all the records have been loaded.iterator()load()Starts the asynchronous loading of records from the pipeline into this dataset.Starts the asynchronous loading of records from the pipeline into this dataset.load(Integer maxRecordsToLoad, JobCallback<DataReader, DataWriter> callback) Starts the asynchronous loading of records from the pipeline into this dataset.setCollectUniqueValues(boolean collectUniqueValues) Indicates if unique values in the dataset should be collected (defaultfalse).protected DatasetsetColumnStatsLoaded(boolean columnStatsLoaded) Sets the flag returned byisColumnStatsLoaded()without callingafterColumnStatsLoaded().setColumnStatsReaderThreads(int columnStatsReaderThreads) The number of threads to use to process column stats (default 2).setDetectBigNumberValues(boolean detectBigNumberValues) Indicates if big decimals and big integers should be analyzed to determine their precision and scale (defaulttrue).setDetectBooleanValues(boolean detectBooleanValues) Indicates if boolean values should be looked for in strings and undefined types (defaulttrue).setDetectNumericValues(boolean detectNumericValues) Indicates if numeric values should be looked for in strings and undefined types (defaulttrue).setDetectTemporalValues(boolean detectTemporalValues) Indicates if date/time patterns should be looked for in strings and undefined types (defaulttrue).setDetectUuidValues(boolean detectUuidValues) Indicates if UUID values should be looked for in strings and undefined types (defaulttrue).setInferStringTypes(boolean inferStringTypes) Indicates if string and undefined types should be analyzed to determine if they represent a numeric, boolean, or date/time value (defaulttrue).setJobExecutor(Executor jobExecutor) setMaxColumnStatsRecords(Long maxColumnStatsRecords) The number of records to use when calculating column stats ornullfor all records (defaultnull).setMaxColumnsToAnalyze(Integer maxColumnsToAnalyze) Defines the first N columns to analyze for stats for column stats.setPipeline(AbstractPipeline pipeline) protected DatasetsetRecordsLoaded(boolean recordsLoaded) Sets the flag returned byisRecordsLoaded()without callingafterRecordsLoaded().stream()Returns a Stream over records cached in this dataset.toRecord()Converts this object's state to a record thatRecordSerializable.fromRecord(Record)can load.protected voidupdateColumns(Record record, DataWriter writer) Queues the record's fields for asynchronous column stats until the writer has writtengetMaxColumnStatsRecords()records; subclass writers call this for each record they cache.Blocks until column stats collection ends, checking once a second; returns at once if none is running.waitForColumnStatsToLoad(long minRecords, long maxWaitTimeMillis) Blocks until column stats coverminRecordsrecords, collection ends ormaxWaitTimeMillismilliseconds pass.Blocks until all records are loaded (seewaitUntilJobFinished()); column stats may still be in progress.waitForRecordsToLoad(long minRecords, long maxWaitTimeMillis) Blocks until at leastminRecordsrecords are loaded, loading ends ormaxWaitTimeMillismilliseconds pass.Blocks until the job returned bygetJob()finishes; returns at once if there is none.Methods inherited from class com.northconcepts.datapipeline.foundations.core.FoundationObject
addExceptionProperties, assertValid, assertValid, clone, exception, exception, exception, getInternalId, getInternalName, resetInternalIdMethods inherited from class java.lang.Object
equals, getClass, hashCode, notify, notifyAll, wait, wait, waitMethods inherited from interface com.northconcepts.datapipeline.core.DataExceptionContributor
addExceptionPropertiesMethods inherited from interface java.lang.Iterable
spliteratorMethods inherited from interface com.northconcepts.datapipeline.core.RecordSerializable
fromJson, fromJson, toJson, toJson, toJsonMethods inherited from interface com.northconcepts.datapipeline.core.XmlSerializable
fromXml, fromXml, fromXmlElement, toXml, toXml, toXml, toXml, toXml, toXmlElement
-
Constructor Details
-
Dataset
Creates a dataset that caches the given pipeline's records onceload()is called.
-
-
Method Details
-
close
public void close()- Specified by:
closein interfaceAutoCloseable- Specified by:
closein interfaceCloseable
-
finalize
-
getPipeline
-
setPipeline
-
getRecordCount
public abstract long getRecordCount()Returns the number of records loaded into this dataset so far. -
getRecord
Returns the loaded record at the given 0-based index. -
createDataReader
Reads all data cached in this Dataset. -
createDataReader
Reads a subset of data cached in this Dataset. -
getRecordList
Get a subset of the records cached in this dataset. -
iterator
-
forEach
Performs the given action for each record cached in this dataset. -
stream
Returns a Stream over records cached in this dataset. -
getColumnCount
public abstract long getColumnCount() -
getColumnNames
-
getColumn
-
getColumn
Returns the column with the given name or null if there is none. -
getOrCreateColumn
Returns the named column, creating it with the given 0-based field index if needed; called while collecting column stats. -
getColumns
-
getMaxColumnStatsRecords
The number of records to use when calculating column stats ornullfor all records (defaultnull). -
setMaxColumnStatsRecords
The number of records to use when calculating column stats ornullfor all records (defaultnull). -
isInferStringTypes
public boolean isInferStringTypes()Indicates if string and undefined types should be analyzed to determine if they represent a numeric, boolean, or date/time value (defaulttrue). -
setInferStringTypes
Indicates if string and undefined types should be analyzed to determine if they represent a numeric, boolean, or date/time value (defaulttrue). -
isDetectTemporalValues
public boolean isDetectTemporalValues()Indicates if date/time patterns should be looked for in strings and undefined types (defaulttrue). -
setDetectTemporalValues
Indicates if date/time patterns should be looked for in strings and undefined types (defaulttrue). -
isDetectNumericValues
public boolean isDetectNumericValues()Indicates if numeric values should be looked for in strings and undefined types (defaulttrue). -
setDetectNumericValues
Indicates if numeric values should be looked for in strings and undefined types (defaulttrue). -
isDetectBooleanValues
public boolean isDetectBooleanValues()Indicates if boolean values should be looked for in strings and undefined types (defaulttrue). -
setDetectBooleanValues
Indicates if boolean values should be looked for in strings and undefined types (defaulttrue). -
isDetectBigNumberValues
public boolean isDetectBigNumberValues()Indicates if big decimals and big integers should be analyzed to determine their precision and scale (defaulttrue). -
setDetectBigNumberValues
Indicates if big decimals and big integers should be analyzed to determine their precision and scale (defaulttrue). -
isDetectUuidValues
public boolean isDetectUuidValues()Indicates if UUID values should be looked for in strings and undefined types (defaulttrue). -
setDetectUuidValues
Indicates if UUID values should be looked for in strings and undefined types (defaulttrue). -
isCollectUniqueValues
public boolean isCollectUniqueValues()Indicates if unique values in the dataset should be collected (defaultfalse). -
setCollectUniqueValues
Indicates if unique values in the dataset should be collected (defaultfalse). -
isDataLoading
public boolean isDataLoading()Returntrueif the records or column stats are currently being loaded. -
isDataLoaded
public boolean isDataLoaded()Returntrueif all the records have been loaded and all the column stats have been loaded. -
isRecordsLoaded
public boolean isRecordsLoaded()Returntrueif all the records have been loaded. The column stats might not have been loaded even when this method returnstruesince they require additional processing.- See Also:
-
setRecordsLoaded
Sets the flag returned byisRecordsLoaded()without callingafterRecordsLoaded(). -
isColumnStatsLoaded
public boolean isColumnStatsLoaded()Returntrueif all the column stats have been loaded. The records would have already been loaded when this method is called since column stats require additional processing. SeeisRecordsLoaded() -
setColumnStatsLoaded
Sets the flag returned byisColumnStatsLoaded()without callingafterColumnStatsLoaded(). -
getMaxColumnsToAnalyze
Defines the first N columns to analyze for stats for column stats. Ifnull(which is default value), all columns are analyzed (defaultnull). -
setMaxColumnsToAnalyze
Defines the first N columns to analyze for stats for column stats. Ifnull(which is default value), all columns are analyzed (defaultnull). -
getDataLoadException
Returns the exception that failed the current or last load or null if there was none. -
getColumnStatsException
Returns the last exception thrown while collecting column stats or null if there was none. -
getJobExecutor
-
setJobExecutor
Sets theExecutorused to run the internalJob. Passing innullwill result in using the default implementation which callsJob.runAsync(). -
createColumnsDataReader
Returns a reader that emits one summary record per column of this dataset (seeDataset.ColumnsDataReader). -
getJob
Returns the job loading this dataset's records or null before the first load and aftercancelLoad(). -
getColumnStatsReaderThreads
public int getColumnStatsReaderThreads()The number of threads to use to process column stats (default 2). -
setColumnStatsReaderThreads
The number of threads to use to process column stats (default 2). -
getMaxRecordsToLoad
The maximum records to load parameter (maxRecordsToLoad) passed to the last call toload(Integer)orload(Integer, JobCallback). This value isnullifload()was called last ornullwas passed to the other load methods. -
beforeLoad
protected abstract void beforeLoad()Called at the start of the data loading process, but before any records or column stats have been loaded. -
afterLoad
protected abstract void afterLoad()Called at the end of the data loading process after all the records and column stats have been loaded. -
afterRecordsLoaded
protected void afterRecordsLoaded()Called during the data loading process after all the records have been loaded. The column stats are unlikely to have been loaded when this method is called since they require additional processing. -
afterColumnStatsLoaded
protected void afterColumnStatsLoaded()Called during the data loading process after all the column stats have been loaded. The records would have already been loaded when this method is called since column stats require additional processing. -
load
Starts the asynchronous loading of records from the pipeline into this dataset. This method returns immediately and does not wait for loading to complete. SeewaitForRecordsToLoad()andwaitForRecordsToLoad(long, long). -
load
Starts the asynchronous loading of records from the pipeline into this dataset. This method returns immediately and does not wait for loading to complete. SeewaitForRecordsToLoad()andwaitForRecordsToLoad(long, long).- Parameters:
maxRecords- the maximum records to load ornullto load all records.
-
load
Starts the asynchronous loading of records from the pipeline into this dataset. This method returns immediately and does not wait for loading to complete. SeewaitForRecordsToLoad()andwaitForRecordsToLoad(long, long).- Parameters:
maxRecordsToLoad- the maximum records to load ornullto load all records.callback- the object to notify as data is being loaded.
-
cancelLoad
public void cancelLoad()Gracefully terminate the asynchronous data loading and column stats calculation, waiting up to 10 seconds in total. -
waitForRecordsToLoad
Blocks until at leastminRecordsrecords are loaded, loading ends ormaxWaitTimeMillismilliseconds pass. -
waitForRecordsToLoad
Blocks until all records are loaded (seewaitUntilJobFinished()); column stats may still be in progress. -
waitUntilJobFinished
Blocks until the job returned bygetJob()finishes; returns at once if there is none. -
waitForColumnStatsToLoad
Blocks until column stats coverminRecordsrecords, collection ends ormaxWaitTimeMillismilliseconds pass. -
waitForColumnStatsToLoad
Blocks until column stats collection ends, checking once a second; returns at once if none is running. -
createDataWriter
Writes records to this dataset's cache after clearing it. -
updateColumns
Queues the record's fields for asynchronous column stats until the writer has writtengetMaxColumnStatsRecords()records; subclass writers call this for each record they cache. -
addField
Adds the field to the stats of its column, creating the column if needed; runs on a column stats worker thread. -
toRecord
Description copied from interface:RecordSerializableConverts this object's state to a record thatRecordSerializable.fromRecord(Record)can load.- Specified by:
toRecordin interfaceRecordSerializable- Overrides:
toRecordin classBean
-
fromRecord
Description copied from interface:RecordSerializableLoads this instance's state from a record and returnsthis(for fluid API call chaining). For fluid API call chaining, the overridden method should change the declared return type to its class.- Specified by:
fromRecordin interfaceRecordSerializable- Overrides:
fromRecordin classBean- Parameters:
source-- Returns:
- this instance.
-