All Implemented Interfaces:
DataExceptionContributor, JsonSerializable, RecordSerializable, XmlSerializable, Closeable, Serializable, AutoCloseable, Iterable<Record>
Direct Known Subclasses:
JdbcDataset, LocalFileDataset, MemoryDataset, MvStoreDataset

public abstract class Dataset extends FoundationObject implements Iterable<Record>, Closeable
The base class for caching records produced by a Pipeline or DataMappingPipeline. This class handles asynchronous loading and column metadata and statistics.
See Also:
  • Constructor Details

    • Dataset

      public Dataset(AbstractPipeline pipeline)
      Creates a dataset that caches the given pipeline's records once load() is called.
  • Method Details

    • close

      public void close()
      Specified by:
      close in interface AutoCloseable
      Specified by:
      close in interface Closeable
    • finalize

      protected void finalize() throws Throwable
      Overrides:
      finalize in class Object
      Throws:
      Throwable
    • getPipeline

      public AbstractPipeline getPipeline()
    • setPipeline

      public Dataset setPipeline(AbstractPipeline pipeline)
    • getRecordCount

      public abstract long getRecordCount()
      Returns the number of records loaded into this dataset so far.
    • getRecord

      public abstract Record getRecord(long index)
      Returns the loaded record at the given 0-based index.
    • createDataReader

      public DataReader createDataReader()
      Reads all data cached in this Dataset.
    • createDataReader

      public DataReader createDataReader(long offset, int count)
      Reads a subset of data cached in this Dataset.
    • getRecordList

      public RecordList getRecordList(long offset, int count)
      Get a subset of the records cached in this dataset.
    • iterator

      public Iterator<Record> iterator()
      Specified by:
      iterator in interface Iterable<Record>
    • forEach

      public void forEach(Consumer<? super Record> consumer)
      Performs the given action for each record cached in this dataset.
      Specified by:
      forEach in interface Iterable<Record>
    • stream

      public Stream<Record> stream()
      Returns a Stream over records cached in this dataset.
    • getColumnCount

      public abstract long getColumnCount()
    • getColumnNames

      public abstract List<String> getColumnNames()
    • getColumn

      public abstract Column getColumn(int index)
    • getColumn

      public abstract Column getColumn(String name)
      Returns the column with the given name or null if there is none.
    • getOrCreateColumn

      protected abstract Column getOrCreateColumn(String name, int index)
      Returns the named column, creating it with the given 0-based field index if needed; called while collecting column stats.
    • getColumns

      public abstract List<Column> getColumns()
    • getMaxColumnStatsRecords

      public Long getMaxColumnStatsRecords()
      The number of records to use when calculating column stats or null for all records (default null).
    • setMaxColumnStatsRecords

      public Dataset setMaxColumnStatsRecords(Long maxColumnStatsRecords)
      The number of records to use when calculating column stats or null for all records (default null).
    • 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 (default true).
    • setInferStringTypes

      public Dataset setInferStringTypes(boolean inferStringTypes)
      Indicates if string and undefined types should be analyzed to determine if they represent a numeric, boolean, or date/time value (default true).
    • isDetectTemporalValues

      public boolean isDetectTemporalValues()
      Indicates if date/time patterns should be looked for in strings and undefined types (default true).
    • setDetectTemporalValues

      public Dataset setDetectTemporalValues(boolean detectTemporalValues)
      Indicates if date/time patterns should be looked for in strings and undefined types (default true).
    • isDetectNumericValues

      public boolean isDetectNumericValues()
      Indicates if numeric values should be looked for in strings and undefined types (default true).
    • setDetectNumericValues

      public Dataset setDetectNumericValues(boolean detectNumericValues)
      Indicates if numeric values should be looked for in strings and undefined types (default true).
    • isDetectBooleanValues

      public boolean isDetectBooleanValues()
      Indicates if boolean values should be looked for in strings and undefined types (default true).
    • setDetectBooleanValues

      public Dataset setDetectBooleanValues(boolean detectBooleanValues)
      Indicates if boolean values should be looked for in strings and undefined types (default true).
    • isDetectBigNumberValues

      public boolean isDetectBigNumberValues()
      Indicates if big decimals and big integers should be analyzed to determine their precision and scale (default true).
    • setDetectBigNumberValues

      public Dataset setDetectBigNumberValues(boolean detectBigNumberValues)
      Indicates if big decimals and big integers should be analyzed to determine their precision and scale (default true).
    • isDetectUuidValues

      public boolean isDetectUuidValues()
      Indicates if UUID values should be looked for in strings and undefined types (default true).
    • setDetectUuidValues

      public Dataset setDetectUuidValues(boolean detectUuidValues)
      Indicates if UUID values should be looked for in strings and undefined types (default true).
    • isCollectUniqueValues

      public boolean isCollectUniqueValues()
      Indicates if unique values in the dataset should be collected (default false).
    • setCollectUniqueValues

      public Dataset setCollectUniqueValues(boolean collectUniqueValues)
      Indicates if unique values in the dataset should be collected (default false).
    • isDataLoading

      public boolean isDataLoading()
      Return true if the records or column stats are currently being loaded.
    • isDataLoaded

      public boolean isDataLoaded()
      Return true if all the records have been loaded and all the column stats have been loaded.
    • isRecordsLoaded

      public boolean isRecordsLoaded()
      Return true if all the records have been loaded. The column stats might not have been loaded even when this method returns true since they require additional processing.
      See Also:
    • setRecordsLoaded

      protected Dataset setRecordsLoaded(boolean recordsLoaded)
      Sets the flag returned by isRecordsLoaded() without calling afterRecordsLoaded().
    • isColumnStatsLoaded

      public boolean isColumnStatsLoaded()
      Return true if 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. See isRecordsLoaded()
    • setColumnStatsLoaded

      protected Dataset setColumnStatsLoaded(boolean columnStatsLoaded)
      Sets the flag returned by isColumnStatsLoaded() without calling afterColumnStatsLoaded().
    • getMaxColumnsToAnalyze

      public Integer getMaxColumnsToAnalyze()
      Defines the first N columns to analyze for stats for column stats. If null (which is default value), all columns are analyzed (default null).
    • setMaxColumnsToAnalyze

      public Dataset setMaxColumnsToAnalyze(Integer maxColumnsToAnalyze)
      Defines the first N columns to analyze for stats for column stats. If null (which is default value), all columns are analyzed (default null).
    • getDataLoadException

      public Throwable getDataLoadException()
      Returns the exception that failed the current or last load or null if there was none.
    • getColumnStatsException

      public Throwable getColumnStatsException()
      Returns the last exception thrown while collecting column stats or null if there was none.
    • getJobExecutor

      public Executor getJobExecutor()
      Returns the Executor used to run the internal Job. The default implementation calls Job.runAsync().
    • setJobExecutor

      public Dataset setJobExecutor(Executor jobExecutor)
      Sets the Executor used to run the internal Job. Passing in null will result in using the default implementation which calls Job.runAsync().
    • createColumnsDataReader

      public DataReader createColumnsDataReader()
      Returns a reader that emits one summary record per column of this dataset (see Dataset.ColumnsDataReader).
    • getJob

      public Job getJob()
      Returns the job loading this dataset's records or null before the first load and after cancelLoad().
    • getColumnStatsReaderThreads

      public int getColumnStatsReaderThreads()
      The number of threads to use to process column stats (default 2).
    • setColumnStatsReaderThreads

      public Dataset setColumnStatsReaderThreads(int columnStatsReaderThreads)
      The number of threads to use to process column stats (default 2).
    • getMaxRecordsToLoad

      public Integer getMaxRecordsToLoad()
      The maximum records to load parameter (maxRecordsToLoad) passed to the last call to load(Integer) or load(Integer, JobCallback). This value is null if load() was called last or null was 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

      public Dataset 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. See waitForRecordsToLoad() and waitForRecordsToLoad(long, long).
    • load

      public Dataset load(Integer maxRecords)
      Starts the asynchronous loading of records from the pipeline into this dataset. This method returns immediately and does not wait for loading to complete. See waitForRecordsToLoad() and waitForRecordsToLoad(long, long).
      Parameters:
      maxRecords - the maximum records to load or null to load all records.
    • load

      public Dataset load(Integer maxRecordsToLoad, JobCallback<DataReader,DataWriter> callback)
      Starts the asynchronous loading of records from the pipeline into this dataset. This method returns immediately and does not wait for loading to complete. See waitForRecordsToLoad() and waitForRecordsToLoad(long, long).
      Parameters:
      maxRecordsToLoad - the maximum records to load or null to 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

      public Dataset waitForRecordsToLoad(long minRecords, long maxWaitTimeMillis)
      Blocks until at least minRecords records are loaded, loading ends or maxWaitTimeMillis milliseconds pass.
    • waitForRecordsToLoad

      public Dataset waitForRecordsToLoad()
      Blocks until all records are loaded (see waitUntilJobFinished()); column stats may still be in progress.
    • waitUntilJobFinished

      public Dataset waitUntilJobFinished()
      Blocks until the job returned by getJob() finishes; returns at once if there is none.
    • waitForColumnStatsToLoad

      public Dataset waitForColumnStatsToLoad(long minRecords, long maxWaitTimeMillis)
      Blocks until column stats cover minRecords records, collection ends or maxWaitTimeMillis milliseconds pass.
    • waitForColumnStatsToLoad

      public Dataset waitForColumnStatsToLoad()
      Blocks until column stats collection ends, checking once a second; returns at once if none is running.
    • createDataWriter

      protected abstract DataWriter createDataWriter()
      Writes records to this dataset's cache after clearing it.
    • updateColumns

      protected void updateColumns(Record record, DataWriter writer)
      Queues the record's fields for asynchronous column stats until the writer has written getMaxColumnStatsRecords() records; subclass writers call this for each record they cache.
    • addField

      protected Column addField(Record record, Field field, int fieldIndex)
      Adds the field to the stats of its column, creating the column if needed; runs on a column stats worker thread.
    • toRecord

      public Record toRecord()
      Description copied from interface: RecordSerializable
      Converts this object's state to a record that RecordSerializable.fromRecord(Record) can load.
      Specified by:
      toRecord in interface RecordSerializable
      Overrides:
      toRecord in class Bean
    • fromRecord

      public DataWriterPipelineOutput fromRecord(Record source)
      Description copied from interface: RecordSerializable
      Loads this instance's state from a record and returns this (for fluid API call chaining). For fluid API call chaining, the overridden method should change the declared return type to its class.
      Specified by:
      fromRecord in interface RecordSerializable
      Overrides:
      fromRecord in class Bean
      Parameters:
      source -
      Returns:
      this instance.