Ravi Anne
03/25/2025, 4:30 PMimport 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())Robert Howell
03/25/2025, 8:47 PM.collect() in append_iceberg_data.
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.Ravi Anne
03/26/2025, 5:00 AMRobert Howell
03/26/2025, 4:29 PMRobert Howell
03/26/2025, 4:30 PMColin Ho
03/26/2025, 5:50 PMwrite_parquet instead of collect.Robert Howell
03/26/2025, 6:00 PM