Hi folks, we have timeseries data and experimenti...
# daft-dev
r
Hi folks, we have timeseries data and experimenting with daft for analytics. The data is being stored in Iceberg table on S3 (cold tier data) we use daft dataframe to materialize the data, and every minute we fetch delta data, based on time and concatenate. This is akin to warm tier implementation to iceberg data. Do you see any better ways to cache the incremental data ? Also, if dataframe grows in size, does it spill to disk ? Attached the sample code for reference CC: @jay
Copy code
import asyncio
import json
import logging
import os
from pathlib import Path

import daft
from apscheduler.events import EVENT_JOB_EXECUTED
from apscheduler.schedulers.asyncio import AsyncIOScheduler
from daft import DataFrame, col
from pyiceberg.catalog.sql import SqlCatalog
from pyiceberg.partitioning import PartitionField, PartitionSpec
from pyiceberg.schema import Schema
from pyiceberg.transforms import IdentityTransform
from pyiceberg.types import DoubleType, IntegerType, NestedField, StringType

logging.basicConfig(
    level=logging.INFO,
    format="%(asctime)s [%(levelname)s] %(name)s: %(message)s",
    handlers=[
        logging.StreamHandler()
    ]
)
logger = logging.getLogger(__name__)


def default_catalog(tmpdir):

    os.makedirs(f'{tmpdir}', exist_ok=True)

    sch = Schema(
        NestedField(field_id=1,name="CUST_ID", field_type=StringType()),
        NestedField(field_id=2,name="START_DATE",field_type=StringType()),
        NestedField(field_id=3,name="END_DATE",field_type=StringType()),
        NestedField(field_id=4,name="TRANS_ID",field_type=StringType()),
        NestedField(field_id=5,name="AMOUNT",field_type=DoubleType()),
        NestedField(field_id=6,name="EXP_TYPE",field_type=StringType()),
        NestedField(field_id=7,name="DAY",field_type=IntegerType()),
        NestedField(field_id=8,name="MONTH",field_type=IntegerType()),
        NestedField(field_id=9,name="YEAR",field_type=IntegerType()),
        NestedField(field_id=10,name="DATE",field_type=StringType()),
    )

    # Define the partition specification
    partition_spec = PartitionSpec(
        PartitionField(
            source_id=8, field_id=1001, transform=IdentityTransform(), name="month",
        )
    )

    # creates the catalog on the local disk
    catalog = SqlCatalog(
        "default",
        **{
            "uri": f"sqlite:///{tmpdir}/pyiceberg_catalog.db",
            "warehouse": f"file://{tmpdir}",
        },
    )
    catalog.create_namespace_if_not_exists("random")
    table = catalog.create_table_if_not_exists("random.transactions",schema=sch,partition_spec=partition_spec,
        properties={
            "format-version": "2",
        })

    return (catalog, table)


(catalog, tableRef) = default_catalog(Path("/tmp/transactionhouse"))


class DataAppender:

    def append_iceberg_data(self):
            transactions = daft.read_csv(str(Path(__file__).parent/"sample_data.csv")).collect()
            logger.info(transactions.limit(2).to_pydict())
            logger.info(f"schema is: {transactions.schema()}")
            result_df = transactions.write_iceberg(tableRef, mode="append")
            logger.info(f"appended {result_df.count_rows()} records to iceberg table")



class DataCache:
    def __init__(self):
        self.cached_records:DataFrame = None
        self.day = 1
        self.fetch()

    def fetch(self):
        self.cached_records = (daft.read_iceberg(tableRef)
                        .where((col("MONTH") == 3) & (col("DAY") == self.day)).collect())
        logger.info(f"cache initialization records:{self.cached_records.count_rows()}")

    async def daft_refresh(self):
        tableRef.refresh()
        self.day= self.day + 1 
        new_records = daft.read_iceberg(tableRef).where((col("MONTH") == 3) & ((col("DAY") == self.day))).collect()
        new_records = new_records.select(*self.cached_records.schema().column_names())
        logger.info(f"new_records for day {self.day}: {new_records.count_rows()}")
        self.cached_records = (self.cached_records.concat(new_records).groupby('TRANS_ID')
                                .agg(#deduplication
                                    col('CUST_ID').any_value().alias('CUST_ID'),
                                    col('START_DATE').any_value().alias('START_DATE'),
                                    col('END_DATE').any_value().alias('END_DATE'),
                                    col('AMOUNT').any_value().alias('AMOUNT'),
                                    col('EXP_TYPE').any_value().alias('EXP_TYPE'),
                                    col('DAY').any_value().alias('DAY'),
                                    col('MONTH').any_value().alias('MONTH'),
                                    col('YEAR').any_value().alias('YEAR'),
                                    col('DATE').max().alias('DATE')
                                    ).collect()
                                )
        logger.info(f"cached_records:{self.cached_records.count_rows()}")
        # logger.info(f"analysis: {self.analyse()}")
    
    # def analyse(self):
    #     # print(transactions.collect_schema())
    #     res = (
    #         self.cached_records 
    #         .groupby("EXP_TYPE", "YEAR", "MONTH")
    #         .agg(col("AMOUNT").mean())
    #         .sort([col("EXP_TYPE"), col("YEAR"), col("MONTH")])
    #         .collect()
    #     )
    #     dict_data = res.to_pydict()
    #     # Convert dictionary to JSON
    #     json_data = json.dumps(dict_data)
    #     return json_data


async def main():

    asyncevent = asyncio.Event()

    scheduler = AsyncIOScheduler()
    cache_instance = DataCache()
    scheduler.add_job(cache_instance.daft_refresh,'interval',id="123444", seconds=10)

    max_execution_count=5
    execution_count =0
    def job_listener(event):
        nonlocal execution_count
        execution_count += 1
        logger.info(f".....Job executed {execution_count} times........")
        if execution_count >= max_execution_count:
            scheduler.remove_job("123444")
            logger.info(f"Job completed {execution_count} executions, and has been removed.")
            asyncevent.set()
            

    scheduler.add_listener(job_listener, mask=EVENT_JOB_EXECUTED)                                                                                                                                                                                                                                                                                                                                                                    
    scheduler.start()

    # Keep the script running to allow the scheduler to execute jobs
    try:
        await asyncevent.wait()
    finally:
        scheduler.shutdown()


if __name__ == "__main__":
    DataAppender().append_iceberg_data()
    asyncio.run(main())
r
Hey Ravi, one thing I notice is that you may not need to use
.collect()
in append_iceberg_data.
Copy code
transactions = daft.read_csv(str(Path(__file__).parent/"sample_data.csv"))
            <http://logger.info|logger.info>(transactions.limit(2).to_pydict())
            <http://logger.info|logger.info>(f"schema is: {transactions.schema()}")
            result_df = transactions.write_iceberg(tableRef, mode="append")
            <http://logger.info|logger.info>(f"appended {result_df.count_rows()} records to iceberg table")
For better ways to cache the data, are you running into any current issues. The sqlite iceberg catalog and a temp dir is pretty good solution for caching results on disk.
r
Hey Conner, The sample code uses sqlite iceberg catalog for illustration purpose, but in the actual case it would be S3 catalog. Data directly gets appended to S3 catalog. To speed up the queries, we are caching the s3 catalog data in Daft dataframe, via periodic refresh. 1. If the data grows in size, does Daft dataframe spills to disk automatically or anything else to be done ? 2. Alternatively, even if we choose to fetch data from S3 catalog and write to sqlite catalog on disk periodically, to speed up the queries, one has to perform Iceberg maintenance activities on two catalogs now. what would be your guidance ?
r
For (1) I'm tagging @Colin Ho as he can better answer this than me 😄 For (2) I think since the disc catalog is just a cache you shouldn't have to perform any maintenance, as any gains would be marginal as to warrant the complexity not worth it. An interesting idea would be to actually cache the files themselves rather than the dataframes. This would simplify things on your end. https://docs.aws.amazon.com/whitepapers/latest/s3-optimizing-performance-best-practices/using-caching-for-frequently-accessed-content.html
You raise an interesting use-case for sure! We don't have an SOP or specific guidance. Generally caching the recently accessed S3 files should provide improvement without having to manage the cache yourself. Hope that helps.
c
Daft right now does not have automatic spilling mechanisms, 😞 . If the intermediate data is growing too large to fit in memory, I think the simplest way right now would be to do a
write_parquet
instead of
collect
.
r
Thanks Colin!