// Parallel import of multiple CSV files
login(`admin,`123456)

// Create database and table
if (existsDatabase("dfs://sh_entrust"))
{
	dropDatabase("dfs://sh_entrust")
}
create database "dfs://sh_entrust" partitioned by VALUE(2022.01.01..2022.01.03), HASH([SYMBOL, 10]), engine='TSDB'

create table "dfs://sh_entrust"."entrust"(
	SecurityID SYMBOL,
	TransactTime TIMESTAMP,
	valOrderNoue INT,
	Price DOUBLE,
	Balance INT,
	OrderBSFlag SYMBOL,
	OrdType SYMBOL,
	OrderIndex INT,
	ChannelNo INT,
	BizIndex INT)
partitioned by TransactTime,SecurityID,
sortColumns = [`SecurityID,`TransactTime]

// Define type conversion function
def transType(mutable memTable)
{
	return memTable.replaceColumn!(`col0,lpad(string(memTable.col0),6,`0)).replaceColumn!(`col1,datetimeParse(string(memTable.col1),"yyyyMMddHHmmssSSS")).replaceColumn!(`col5,string(memTable.col5)).replaceColumn!(`col6,string(memTable.col6))
}

// Define function to load one day's data
def loadOneDayFile(db,table,filePath)
{
	csvFiles = exec filename from files(filePath) where filename like"%.csv"
	for(csvIdx in csvFiles)
	{
		loadTextEx(dbHandle = db, tableName = table, partitionColumns = `col1`col0, filename = filePath + "/"  + csvIdx, transform = transType, skipRows = 1)
	}
}

// Define function to submit parallel jobs
def parallelLoad(allFileContents)
{
	db = database("dfs://sh_entrust")
	table = `entrust
	dateFiles = exec filename from files(allFileContents) where isDir = true
	for(dateIdx in dateFiles)
	{
		submitJob("parallelLoad" + dateIdx,"parallelLoad",loadOneDayFile{db,table,},allFileContents + "/" + dateIdx)
	}
}

// Call the function, submit jobs — adjust the directory path as needed
allFileContents = "/home/ychan/data/loadForPoc/SH/Order"
parallelLoad(allFileContents)

// Check job execution status
getRecentJobs()