otp.ReadFromIceberg#
- class ReadFromIceberg(catalog_config=None, table_identifier=None, fields='', drop_fields=False, as_of_time=None, time_assignment='end', where='', symbology='', symbol_name_field='SYMBOL_NAME', batch_size=1000, symbol=utils.adaptive, db=utils.adaptive_to_default, tick_type=utils.adaptive, start=utils.adaptive, end=utils.adaptive, schema=None, query_parameters=None, **kwargs)#
Bases:
SourceRead ticks from an Iceberg table.
The table is specified by
table_identifierand is accessed with the catalog settings from thecatalog_configfile.- Parameters:
catalog_config (str) – Path to a configuration file containing the Iceberg catalog parameters. For details on the file structure and how it is processed, see Examples below.
table_identifier (str) – The Iceberg table identifier, including the namespace and table name (for example,
namespace.tablename).fields (list, str) – A list of fields (
listor comma-separated string) to be included in the output ticks. If not set, all available fields from the table are returned. Seedrop_fieldsparameter to propagate all fields except the listed ones.drop_fields (bool) – If set to
True, the fields listed in thefieldsparameter will not be propagated, and the fields that are not listed there will be propagated instead.as_of_time (
otp.datetime,datetime.datetime, int) – Timestamp for the Iceberg time travel query. The table snapshot that was current as of this timestamp will be read, following the standard Iceberg time travel semantics. Integers are treated as the number of milliseconds since epoch. Datetime values without timezone are treated as values inotp.config.tz.time_assignment (str) – Timestamps of the ticks created by
ReadFromIcebergare set to the start or to the end of the query depending on thetime_assignmentparameter. Possible values arestartandend(for_START_TIMEand_END_TIME).where (str) – Specifies a criterion for selecting the rows to propagate. Wherever possible, Iceberg database predicate push-down is performed for certain sub-clauses of the specified expression.
symbology (str) – Symbology of the symbol names found in the
symbol_name_fieldfield. If specified, the reference database is used to find synonyms for the query symbols within this symbology.symbol_name_field (str) – Field that is expected to contain the symbol name. When this parameter is set and specific symbols are bound to the query, only rows matching those symbols are propagated. If multiple symbols are present, rows are routed to the appropriate sub-queries based on this field’s value.
batch_size (int) – The maximum number of rows to be collected in the Java layer before they are transferred to the C++ layer.
symbol (str, list of str,
Source,query,eval query) – Symbol(s) from which data should be taken.tick_type (str) – Tick type. Default: ANY.
start (
otp.datetime) – Start time for tick generation. By default the start time of the query will be used.end (
otp.datetime) – End time for tick generation. By default the end time of the query will be used.schema (dict) –
Set the schema of the python
Sourceobject of this class.Schema can’t be automatically derived from the Iceberg table, so it should be set manually for Python-level type checking to work.
query_parameters (
otp.QueryParameters) – Additional query properties to be set in the resulting .otq file. They will be used if they are not overridden by other parameters or inotp.run.kwargs – Deprecated. Use
schemainstead. Dictionary of columns names with their types.
Note
This EP is supported only on OneTick for 64-bit Windows/Linux platforms.
This source requires Java 17 or higher to be installed. Also Java library path should be added to PATH (Windows) or LD_LIBRARY_PATH (Linux) environment variables and to OMD_JAVALIBPATH OneTick config variable.
Examples
Write Iceberg catalog configuration first:
>>> with open('catalog.cfg', 'w') as f: ... f.write(''' ... catalog.name=rest_catalog ... type=rest ... uri=http://localhost:8181 ... warehouse=s3://warehouse/ ... s3.region=us-east-1 ... client.region=us-east-1 ... s3.path-style-access=true ... s3.access-key-id=admin ... s3.secret-access-key=password ... io-impl=org.apache.iceberg.aws.s3.S3FileIO ... s3.delete.enabled=true ... s3.acceleration-enabled=false ... ''')
Then read all fields from the table:
>>> data = otp.ReadFromIceberg(catalog_config='catalog.cfg', ... table_identifier='exchange.trades') >>> otp.run(data)
Read only the fields you need:
>>> data = otp.ReadFromIceberg(catalog_config='catalog.cfg', ... table_identifier='exchange.trades', ... fields=['PRICE', 'SIZE']) >>> otp.run(data)
Propagate all fields except the listed ones:
>>> data = otp.ReadFromIceberg(catalog_config='catalog.cfg', ... table_identifier='exchange.trades', ... fields=['SIZE'], ... drop_fields=True) >>> otp.run(data)
Filter rows and set tick timestamps from the
trade_timefield:>>> data = otp.ReadFromIceberg(catalog_config='catalog.cfg', ... table_identifier='exchange.trades', ... fields=['PRICE', 'SIZE', 'trade_time'], ... time_assignment='trade_time', ... where='PRICE > 20') >>> otp.run(data)
Read the table snapshot that was current at the specified point in time:
>>> data = otp.ReadFromIceberg(catalog_config='catalog.cfg', ... table_identifier='exchange.trades', ... fields=['PRICE', 'SIZE'], ... as_of_time=otp.datetime(2023, 1, 1, 12)) >>> otp.run(data)
See also
READ_FROM_ICEBERG OneTick event processor