【问题标题】:Buying ram to avoid chunking for 30-50Gb plus files购买 ram 以避免 30-50Gb plus 文件的分块
【发布时间】:2016-11-18 13:19:19
【问题描述】:

我使用 pandas 读取非常大的 csv 文件,这些文件也是 gzip 压缩的。 我解压缩成大约 30-50GB 的 csv 文件。 我将文件分块并处理/操作它们。 最后将相关数据添加到我压缩的 HDF5 文件中

它工作得很好,但速度很慢,因为我每天必须处理一个文件并且有几年的数据(600TB 未压缩的 csv)

购买更多的内存是避免分块和加快进程(例如 64GB/128GB)的好方法吗? 但这会让 pandas 变得又慢又笨重吗? 我是否正确地说切换到 C++ 可以加快这个过程,但我仍然受到读取过程的影响,并且不得不分块处理数据。 最后,有没有人对处理这个问题的最佳方法有任何想法。

顺便说一句,一旦工作完成,我就不必返回并再次处理数据,所以只想让它在合理的时间内工作,所以写一些并行过程可能很好但经验有限的东西我需要一些时间来构建它,所以除非这是唯一的选择,否则我宁愿不要。

更新。我认为查看代码会更容易。无论如何,我不相信代码特别慢。我认为技术/方法可能是。

def txttohdf(path, contract):
    #create dataframes for trade and quote
    dftrade = pd.DataFrame(columns = ["datetime", "Price", "Volume"])
    dfquote = pd.DataFrame(columns = ["datetime", "BidPrice", "BidSize","AskPrice", "AskSize"])
    #create an hdf5 file with high compression and table so we can append
    hdf = pd.HDFStore(path + contract + '.h5', complevel=9, complib='blosc')
    hdf.put('trade', dftrade, format='table', data_columns=True)
    hdf.put('quote', dfquote, format='table', data_columns=True)
    #date1 = date(start).strftime('%Y%m%d')
    #date2 = date(end).strftime('%Y%m%d')
    #dd = [date1 + timedelta(days=x) for x in range((date2-date1).days + 1)]
    #walkthrough directories
    for subdir, dir, files in os.walk(path):
        for file in files:
            #check if contract has name
            #print(file)
                #create filename from directory and file 

            filename = os.path.join(subdir, file)
                #read in csv
            if filename.endswith('.gz'):

                df = pd.read_csv(gzip.open(filename),header=0,iterator=True,chunksize = 10000, low_memory =False,  names = ['RIC','Date','Time','GMTOffset','Type','ExCntrbID','LOC','Price','Volume','MarketVWAP','BuyerID','BidPrice','BidSize','NoBuyers','SellerID','AskPrice','AskSize','NoSellers','Qualifiers','SeqNo','ExchTime','BlockTrd','FloorTrd','PERatio','Yield','NewPrice','NewVol','NewSeqNo','BidYld','AskYld','ISMABidYld','ISMAAskYld','Duration','ModDurtn','BPV','AccInt','Convexity','BenchSpd','SwpSpd','AsstSwpSpd','SwapPoint','BasePrice','UpLimPrice','LoLimPrice','TheoPrice','StockPrice','ConvParity','Premium','BidImpVol','AskImpVol','ImpVol','PrimAct','SecAct','GenVal1','GenVal2','GenVal3','GenVal4','GenVal5','Crack','Top','FreightPr','1MnPft','3MnPft','PrYrPft','1YrPft','3YrPft','5YrPft','10YrPft','Repurch','Offer','Kest','CapGain','Actual','Prior','Revised','Forecast','FrcstHigh','FrcstLow','NoFrcts','TrdQteDate','QuoteTime','BidTic','TickDir','DivCode','AdjClose','PrcTTEFlag','IrgTTEFlag','PrcSubMktId','IrgSubMktId','FinStatus','DivExDate','DivPayDate','DivAmt','Open','High','Low','Last','OpenYld','HighYld','LowYld','ShortPrice','ShortVol','ShortTrdVol','ShortTurnnover','ShortWeighting','ShortLimit','AccVolume','Turnover','ImputedCls','ChangeType','OldValue','NewValue','Volatility','Strike','Premium','AucPrice','Auc Vol','MidPrice','FinEvalPrice','ProvEvalPrice','AdvancingIssues','DecliningIssues','UnchangedIssues','TotalIssues','AdvancingVolume','DecliningVolume','UnchangedVolume','TotalVolume','NewHighs','NewLows','TotalMoves','PercentageChange','AdvancingMoves','DecliningMoves','UnchangedMoves','StrongMarket','WeakMarket','ChangedMarket','MarketVolatility','OriginalDate','LoanAskVolume','LoanAskAmountTradingPrice','PercentageShortVolumeTradedVolume','PercentageShortPriceTradedPrice','ForecastNAV','PreviousDaysNAV','FinalNAV','30DayATMIVCall','60DayATMIVCall','90DayATMIVCall','30DayATMIVPut','60DayATMIVPut','90DayATMIVPut','BackgroundReference','DataSource','BidSpread','AskSpread','ContractPhysicalUnits','Miniumumquantity','NumberPhysicals','ClosingReferencePrice','ImbalanceQuantity','FarClearingPrice','NearClearingPrice','OptionAdjustedSpread','ZSpread','ConvexityPremium','ConvexityRatio','PercentageDailyReturn','InterpolatedCDSBasis','InterpolatedCDSSpread','ClosesttoMaturityCDSBasis','SettlementDate','EquityPrice','Parity','CreditSpread','Delta','InputVolatility','ImpliedVolatility','FairPrice','BondFloor','Edge','YTW','YTB','SimpleMargin','DiscountMargin','12MonthsEPS','UpperTradingLimit','LowerTradingLimit','AmountOutstanding','IssuePrice','GSpread','MiscValue','MiscValueDescription'])
                #parse date time this is quicker than doing it while we read it in
                for chunk in df:
                    chunk['datetime'] = chunk.apply(lambda row: datetime.datetime.strptime(row['Date']+ ':' + row['Time'],'%d-%b-%Y:%H:%M:%S.%f'), axis=1)
                    #df = df[~df.comment.str.contains('ALIAS')]
                #drop uneeded columns inc date and time
                    chunk = chunk.drop(['Date','Time','GMTOffset','ExCntrbID','LOC','MarketVWAP','BuyerID','NoBuyers','SellerID','NoSellers','Qualifiers','SeqNo','ExchTime','BlockTrd','FloorTrd','PERatio','Yield','NewPrice','NewVol','NewSeqNo','BidYld','AskYld','ISMABidYld','ISMAAskYld','Duration','ModDurtn','BPV','AccInt','Convexity','BenchSpd','SwpSpd','AsstSwpSpd','SwapPoint','BasePrice','UpLimPrice','LoLimPrice','TheoPrice','StockPrice','ConvParity','Premium','BidImpVol','AskImpVol','ImpVol','PrimAct','SecAct','GenVal1','GenVal2','GenVal3','GenVal4','GenVal5','Crack','Top','FreightPr','1MnPft','3MnPft','PrYrPft','1YrPft','3YrPft','5YrPft','10YrPft','Repurch','Offer','Kest','CapGain','Actual','Prior','Revised','Forecast','FrcstHigh','FrcstLow','NoFrcts','TrdQteDate','QuoteTime','BidTic','TickDir','DivCode','AdjClose','PrcTTEFlag','IrgTTEFlag','PrcSubMktId','IrgSubMktId','FinStatus','DivExDate','DivPayDate','DivAmt','Open','High','Low','Last','OpenYld','HighYld','LowYld','ShortPrice','ShortVol','ShortTrdVol','ShortTurnnover','ShortWeighting','ShortLimit','AccVolume','Turnover','ImputedCls','ChangeType','OldValue','NewValue','Volatility','Strike','Premium','AucPrice','Auc Vol','MidPrice','FinEvalPrice','ProvEvalPrice','AdvancingIssues','DecliningIssues','UnchangedIssues','TotalIssues','AdvancingVolume','DecliningVolume','UnchangedVolume','TotalVolume','NewHighs','NewLows','TotalMoves','PercentageChange','AdvancingMoves','DecliningMoves','UnchangedMoves','StrongMarket','WeakMarket','ChangedMarket','MarketVolatility','OriginalDate','LoanAskVolume','LoanAskAmountTradingPrice','PercentageShortVolumeTradedVolume','PercentageShortPriceTradedPrice','ForecastNAV','PreviousDaysNAV','FinalNAV','30DayATMIVCall','60DayATMIVCall','90DayATMIVCall','30DayATMIVPut','60DayATMIVPut','90DayATMIVPut','BackgroundReference','DataSource','BidSpread','AskSpread','ContractPhysicalUnits','Miniumumquantity','NumberPhysicals','ClosingReferencePrice','ImbalanceQuantity','FarClearingPrice','NearClearingPrice','OptionAdjustedSpread','ZSpread','ConvexityPremium','ConvexityRatio','PercentageDailyReturn','InterpolatedCDSBasis','InterpolatedCDSSpread','ClosesttoMaturityCDSBasis','SettlementDate','EquityPrice','Parity','CreditSpread','Delta','InputVolatility','ImpliedVolatility','FairPrice','BondFloor','Edge','YTW','YTB','SimpleMargin','DiscountMargin','12MonthsEPS','UpperTradingLimit','LowerTradingLimit','AmountOutstanding','IssuePrice','GSpread','MiscValue','MiscValueDescription'], axis=1)
                # convert to datetime explicitly and add nanoseconds to same time stamps
                    chunk['datetime'] = pd.to_datetime(chunk.datetime)
                #nanoseconds = df.groupby(['datetime']).cumcount()
                #df['datetime'] += np.array(nanoseconds, dtype='m8[ns]')  
                # drop empty prints and make sure all prices are valid
                    dfRic = chunk[(chunk["RIC"] == contract)]
                    if len(dfRic)>0:
                        print(dfRic)
                    if ~chunk.empty:
                        dft = dfRic[(dfRic["Type"] == "Trade")]
                        dft.dropna(subset = ["Volume"], inplace =True)
                        dft = dft.drop(["RIC","Type","BidPrice", "BidSize", "AskPrice", "AskSize"], axis=1)
                        dft = dft[(dft["Price"] > 0)]

                    # clean up bid and ask
                        dfq = dfRic[(dfRic["Type"] == "Quote")]
                        dfq.dropna(how = 'all', subset = ["BidSize","AskSize"], inplace =True)
                        dfq = dfq.drop(["RIC","Type","Price", "Volume"], axis=1)
                        dfq = dfq[(dfq["BidSize"] > 0) | (dfq["AskSize"] > 0)]
                        dfq = dfq.ffill()
                    else:
                        print("Empty")    
    #add to hdf and close if loop finished
                    hdf.append('trade', dft, format='table', data_columns=True)
                    hdf.append('quote', dfq, format='table', data_columns=True)
    hdf.close()

【问题讨论】:

  • 你能解释一下什么是慢,为什么它是慢的?如果没有更多细节,很难猜测什么会有助于加快这一进程。
  • 您应该尝试分析和测量程序的性能,以确定哪些点是最慢的,以及内存或 CPU 功率是否是限制因素。这将帮助您缩小可以帮助您的特定更改的范围。然后,您还可以将源代码中最慢的部分上传到codereview.stackexchange.com 上的问题中,并寻求有关提高其性能的建议。
  • 我会尝试分块读取压缩后的 CSV,而不是先解压缩它们——这样你应该有更少的 IO(通常是最慢的部分之一)。除此之外,拥有更多 RAM 应该可以让你拥有更大的块,或者如果你的 RAM 大约是与生成的 DF 相比,大两倍。由于开销,同一台服务器/计算机上的并行处理(如果您的意思是 DASK)会使一切变得更糟。如果您需要真正的强大功能,请查看 Apache PySpark SQL,但这意味着对 Hadoop 集群的投资更高——只需我的 2 美分……
  • 您好,贴出代码是为了给您一个更好的主意。我很确定是阅读和分块部分很耗时。
  • 我认为在创建数据帧时使用多处理和指定数据类型会减少时间。

标签: python csv pandas ram chunking


【解决方案1】:

我认为你有很多东西可以优化:

  • 首先只读取您真正需要的列,而不是读取然后删除它们 - 使用 usecols=list_of_needed_columns 参数

  • 增加你的块大小 - 尝试不同的值 - 我会从10**5开始

  • 不要使用 chunk.apply(...) 来转换您的日期时间 - 它非常慢 - 请改用 pd.to_datetime(column, format='...')

  • 您可以在组合多个条件时更有效地过滤数据,而不是一步一步地进行:

【讨论】:

  • 那太好了,将应用更改-您对我忘记的应用绝对正确。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2013-04-17
  • 2023-01-16
  • 1970-01-01
  • 1970-01-01
  • 2015-09-10
相关资源
最近更新 更多