# 巨大なデータの取り扱い
巨大といってもここでは数千万～数億レコード（数GB～100GB未満ぐらい）のデータをPythonで処理するときについて  
色々試した結果を記す。  

In [1]:
import polars as pl
import numpy as np

# 文字列カラムの表示文字数を50文字に設定
pl.Config.set_fmt_str_lengths(50)

polars.config.Config

## ● 約5000万行 × 5列 × 1ファイル（外付けUSBに保存）の場合  
### 条件
* USBのスペック： 
    * USB3.2 Gen1  
    * 容量： 32GB
    * 他の詳細は忘れた

USBに保存してある大きめのファイルについて処理したいとき。  
データの全カラム・全レコードが解析に必要なケースは多くないと思われるので、  
LazyFrameを活用して高速化、メモリ節約を意識する。  
collect()の際はstreaming=Trueにすることで一気にメモリに読み込まずにバッチ的に読み込んで処理する。

In [2]:
! wc -l ../../../sample_data/from_HDD/00_study/sample_data/sample_big_data_1.csv

49061300 ../../../sample_data/from_HDD/big_data/sample_big_data_1.csv


In [6]:
# 読み込み条件定義
input_file = '../../../sample_data/from_HDD/00_Study/sample_data/sample_big_data_1.csv'
col_names_dtypes = {
    'datetime_col': pl.Utf8, 
    'value_1': pl.Float32, 
    'value_2': pl.Float32, 
    'value_3': pl.Float32, 
    'labels': pl.Utf8}

In [7]:
# 最初の５行だけ読み込み
(
    pl.scan_csv(input_file, has_header=False, schema=col_names_dtypes)
    .head().collect()
)

datetime_col,value_1,value_2,value_3,labels
str,f32,f32,f32,str
"""2023-11-03 13:39:36.088815""",404.69751,443.679321,87.651566,"""CCCC"""
"""2023-11-03 13:39:36.088875""",457.07785,634.807251,549.272156,"""EEEE"""
"""2023-11-03 13:39:36.088891""",620.065186,173.628784,219.745941,"""BBBB"""
"""2023-11-03 13:39:36.088906""",25.579391,143.887222,408.059845,"""BBBB"""
"""2023-11-03 13:39:36.088920""",654.784729,528.226379,1026.964355,"""BBBB"""


最初の5行を表示したりは全く問題ない。

In [7]:
%%time
# filter,with_columnsをやってみる。
(
    pl.scan_csv(input_file, has_header=False, schema=col_names_dtypes, truncate_ragged_lines=True)
    .filter(pl.col('labels') == 'EEEE')
    .with_columns(
        val_1_90tile = pl.col('value_1').quantile(0.9)
    ).collect(streaming=True)
)

CPU times: user 35.8 s, sys: 2min 41s, total: 3min 17s
Wall time: 1min 6s


datetime_col,value_1,value_2,value_3,labels,val_1_90tile
str,f32,f32,f32,str,f32
"""2023-11-03 13:39:36.088875""",457.07785,634.807251,549.272156,"""EEEE""",900.028076
"""2023-11-03 13:39:36.088961""",484.879608,265.178833,711.663208,"""EEEE""",900.028076
"""2023-11-03 13:39:36.089092""",713.928345,101.049217,836.95105,"""EEEE""",900.028076
"""2023-11-03 13:39:36.089165""",964.206848,974.932983,227.194,"""EEEE""",900.028076
"""2023-11-03 13:39:36.089313""",864.268555,748.114075,447.793671,"""EEEE""",900.028076
"""2023-11-03 13:39:36.089325""",730.194397,89.366402,932.182983,"""EEEE""",900.028076
"""2023-11-03 13:39:36.089375""",140.499466,857.064331,237.348999,"""EEEE""",900.028076
"""2023-11-03 13:39:36.089387""",719.163879,674.034851,836.932312,"""EEEE""",900.028076
"""2023-11-03 13:39:36.089461""",832.301819,503.263397,335.356476,"""EEEE""",900.028076
"""2023-11-03 13:39:36.089485""",359.132294,8.830131,935.493713,"""EEEE""",900.028076


filterしてquantileを計算するのは1分ほどかかった。

In [10]:
%%time
# filter,with_columnsをやってみる。
(
    pl.scan_csv(input_file, has_header=False, schema=col_names_dtypes)
    .filter(pl.col('labels') == 'EEEE')
    .select(pl.col('value_1').quantile(0.9))
    .collect(streaming=True)
)

CPU times: user 21.1 s, sys: 3min 1s, total: 3min 22s
Wall time: 53.8 s


value_1
f32
900.028076


In [11]:
%%time
# streamingでfilter,with_columnsをやってみる。
(
    pl.scan_csv(input_file, has_header=False, schema=col_names_dtypes)
    .filter(pl.col('labels') == 'EEEE')
    .collect(streaming=True)
)

CPU times: user 32.8 s, sys: 2min 47s, total: 3min 20s
Wall time: 1min 6s


datetime_col,value_1,value_2,value_3,labels
str,f32,f32,f32,str
"""2023-11-03 13:39:36.088875""",457.07785,634.807251,549.272156,"""EEEE"""
"""2023-11-03 13:39:36.088961""",484.879608,265.178833,711.663208,"""EEEE"""
"""2023-11-03 13:39:36.089092""",713.928345,101.049217,836.95105,"""EEEE"""
"""2023-11-03 13:39:36.089165""",964.206848,974.932983,227.194,"""EEEE"""
"""2023-11-03 13:39:36.089313""",864.268555,748.114075,447.793671,"""EEEE"""
"""2023-11-03 13:39:36.089325""",730.194397,89.366402,932.182983,"""EEEE"""
"""2023-11-03 13:39:36.089375""",140.499466,857.064331,237.348999,"""EEEE"""
"""2023-11-03 13:39:36.089387""",719.163879,674.034851,836.932312,"""EEEE"""
"""2023-11-03 13:39:36.089461""",832.301819,503.263397,335.356476,"""EEEE"""
"""2023-11-03 13:39:36.089485""",359.132294,8.830131,935.493713,"""EEEE"""


### ● selectしてquantileだけ求める場合

In [12]:
%%time
# filter,with_columnsをやってみる。
(
    pl.scan_csv(input_file, has_header=False, schema=col_names_dtypes)
    .select(pl.col('value_1').quantile(0.9))
    .collect(streaming=True)
)

CPU times: user 34.1 s, sys: 3min 4s, total: 3min 38s
Wall time: 51 s


value_1
f32
900.036194


In [None]:
%%time
# filter,with_columnsをやってみる。
(
    pl.scan_csv(input_file, has_header=False, schema=col_names_dtypes)
    .select([
        pl.col('value_1').quantile(0.9),
        pl.col('value_2').quantile(0.9),
        pl.col('value_3').quantile(0.9),
    ])
    .collect(streaming=True)
)

CPU times: user 1min 19s, sys: 2min 50s, total: 4min 10s
Wall time: 1min 3s


value_1,value_2,value_3
f32,f32,f32
900.036194,900.887268,989.969238


quantileの計算を増やしても時間はそれほど変わらなかった。

In [18]:
%%time
# filter,with_columnsをやってみる。
(
    pl.scan_csv(input_file, has_header=False, schema=col_names_dtypes)
    .select([
        pl.col('value_1').quantile(i).alias(f'q_{i*100:.0f}') for i in np.linspace(0.0,1.0,10)
    ])
    .collect(streaming=True)
)

CPU times: user 3min 18s, sys: 3min 1s, total: 6min 19s
Wall time: 1min 4s


q_0,q_11,q_22,q_33,q_44,q_56,q_67,q_78,q_89,q_100
f32,f32,f32,f32,f32,f32,f32,f32,f32,f32
5.5e-05,111.199875,222.284958,333.375244,444.49234,555.536194,666.622375,777.792542,888.939148,1000.0


----------------------------------------------

----------------------------------------------

----------------------------------------------

## ● 約5000万行 × 5列 × 1ファイル（内部SSDに保存）の場合  
### 条件
* 内部SSDのスペック
    * 容量：2TB
    * その他スペック忘れた

内部SSDをWSL2にマウントしてそこのパスを指定して処理する場合。

In [None]:
! wc -l ../../../sample_data/from_HDD/00_study/sample_data/sample_big_data_1.csv

49061300 ../../../sample_data/from_HDD/big_data/sample_big_data_1.csv


In [8]:
# 読み込み条件定義
input_file = '../../../sample_data/from_HDD/00_Study/sample_data/sample_big_data_1.csv'
col_names_dtypes = {
    'datetime_col': pl.Utf8, 
    'value_1': pl.Float32, 
    'value_2': pl.Float32, 
    'value_3': pl.Float32, 
    'labels': pl.Utf8}

In [9]:
(
    pl.scan_csv(input_file, has_header=False, schema=col_names_dtypes)
    .head().collect()
)

datetime_col,value_1,value_2,value_3,labels
str,f32,f32,f32,str
"""2023-11-03 13:39:36.088815""",404.69751,443.679321,87.651566,"""CCCC"""
"""2023-11-03 13:39:36.088875""",457.07785,634.807251,549.272156,"""EEEE"""
"""2023-11-03 13:39:36.088891""",620.065186,173.628784,219.745941,"""BBBB"""
"""2023-11-03 13:39:36.088906""",25.579391,143.887222,408.059845,"""BBBB"""
"""2023-11-03 13:39:36.088920""",654.784729,528.226379,1026.964355,"""BBBB"""


In [10]:
(
    pl.scan_csv(input_file, has_header=False, schema=col_names_dtypes)
    .select(
        pl.col('value_1').mean()
    ).collect(streaming=True)
)

value_1
f32
500.024323


In [17]:
%%time
# filter,with_columnsをやってみる。
(
    pl.scan_csv(input_file, has_header=False, schema=col_names_dtypes)
    .select([
        pl.col('value_1').quantile(i).alias(f'q_{i*100:.0f}') for i in np.linspace(0.01, 1.0, 100)
    ])
    .collect(streaming=True).transpose()
)

CPU times: user 35min 53s, sys: 5min 7s, total: 41min 1s
Wall time: 4min 13s


column_0
f32
10.025609
20.031404
30.018518
40.021393
50.027016
60.030224
70.055244
80.050735
90.06752
100.069817


内部SSDを指定して1%ile～100%ileの計算を実行する場合、上記の通り4分ほどかかった。  

----------------------------------------------

----------------------------------------------

----------------------------------------------

## ● 約5000万行 × 5列 × 1ファイル（WSL2管理下に保存）の場合  
### 条件
* ー

WSL2管理下に大容量ファイルを持ってきて処理する場合。  
ディスクは食ってしまうが、これが一番早い。  
（ただし、サーバで出力されたデータの場合はこの環境まで持ってくるのにも時間がかかることは留意すること。）

In [2]:
# 読み込み条件定義
input_file = '../../../sample_data/Big_data_sample/sample_big_data_1.csv'
col_names_dtypes = {
    'datetime_col': pl.Utf8, 
    'value_1': pl.Float32, 
    'value_2': pl.Float32, 
    'value_3': pl.Float32, 
    'labels': pl.Utf8}

In [3]:
(
    pl.scan_csv(input_file, has_header=False, schema=col_names_dtypes)
    .head().collect()
)

datetime_col,value_1,value_2,value_3,labels
str,f32,f32,f32,str
"""2023-11-03 13:39:36.088815""",404.69751,443.679321,87.651566,"""CCCC"""
"""2023-11-03 13:39:36.088875""",457.07785,634.807251,549.272156,"""EEEE"""
"""2023-11-03 13:39:36.088891""",620.065186,173.628784,219.745941,"""BBBB"""
"""2023-11-03 13:39:36.088906""",25.579391,143.887222,408.059845,"""BBBB"""
"""2023-11-03 13:39:36.088920""",654.784729,528.226379,1026.964355,"""BBBB"""


In [5]:
(
    pl.scan_csv(input_file, has_header=False, schema=col_names_dtypes)
    .select(
        pl.col('value_1').mean()
    ).collect(streaming=True)
)

value_1
f32
500.024323


In [6]:
! time awk -F, '{m+=$2} END{print m/NR;}' ../../../sample_data/Big_data_sample/sample_big_data_1.csv

500.024

real	0m24.167s
user	0m22.612s
sys	0m1.551s


外付けHDDや内部SSDから読み込む場合に比べて、平均の計算にかかる時間が1/15程度になった。  
またawkで平均を求めるよりも早い。  
（上記はWSL2上のコンテナからawkを実行した場合。WSL2から直接実行するとなぜか1分以上かかった）

In [18]:
%%time
# filter,with_columnsをやってみる。
(
    pl.scan_csv(input_file, has_header=False, schema=col_names_dtypes)
    .select([
        pl.col('value_1').quantile(i).alias(f'q_{i*100:.1f}') for i in np.linspace(0.1, 1.0, 100)
    ])
    .collect(streaming=True)
)

CPU times: user 30min 18s, sys: 7.67 s, total: 30min 25s
Wall time: 2min 32s


q_10.0,q_10.9,q_11.8,q_12.7,q_13.6,q_14.5,q_15.5,q_16.4,q_17.3,q_18.2,q_19.1,q_20.0,q_20.9,q_21.8,q_22.7,q_23.6,q_24.5,q_25.5,q_26.4,q_27.3,q_28.2,q_29.1,q_30.0,q_30.9,q_31.8,q_32.7,q_33.6,q_34.5,q_35.5,q_36.4,q_37.3,q_38.2,q_39.1,q_40.0,q_40.9,q_41.8,q_42.7,…,q_67.3,q_68.2,q_69.1,q_70.0,q_70.9,q_71.8,q_72.7,q_73.6,q_74.5,q_75.5,q_76.4,q_77.3,q_78.2,q_79.1,q_80.0,q_80.9,q_81.8,q_82.7,q_83.6,q_84.5,q_85.5,q_86.4,q_87.3,q_88.2,q_89.1,q_90.0,q_90.9,q_91.8,q_92.7,q_93.6,q_94.5,q_95.5,q_96.4,q_97.3,q_98.2,q_99.1,q_100.0
f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,…,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32,f32
100.069817,109.174461,118.256416,127.342415,136.4328,145.513031,154.58725,163.692505,172.790924,181.872528,190.953415,200.050797,209.163757,218.256317,227.341736,236.411926,245.513458,254.605148,263.701355,272.815277,281.867584,290.943298,300.030121,309.125885,318.221619,327.301422,336.406128,345.50354,354.599609,363.690155,372.785736,381.860046,390.947632,400.063232,409.146454,418.236603,427.338165,…,672.676392,681.774963,690.872375,699.943115,709.051697,718.145081,727.232422,736.329529,745.434448,754.543945,763.627075,772.741638,781.838379,790.9021,800.007996,809.112793,818.20752,827.299316,836.39978,845.494995,854.611877,863.676453,872.771545,881.860229,890.955078,900.036133,909.107849,918.206299,927.302917,936.388123,945.502502,954.601013,963.686218,972.767944,981.836914,990.915833,1000.0


パーセンタイルを大量に計算する場合についても40%程早くなった。

In [20]:
%%time
df = (
    pl.scan_csv(input_file, has_header=False, schema=col_names_dtypes)
)

col_1 = df.select([pl.col('value_1').quantile(i).alias(f'q_{i*100:.1f}') for i in np.linspace(0.1, 1.0, 10)]).collect(streaming=True)
col_2 = df.select([pl.col('value_2').quantile(i).alias(f'q_{i*100:.1f}') for i in np.linspace(0.1, 1.0, 10)]).collect(streaming=True)

CPU times: user 6min 35s, sys: 5.06 s, total: 6min 40s
Wall time: 34.6 s


In [21]:
col_1

q_10.0,q_20.0,q_30.0,q_40.0,q_50.0,q_60.0,q_70.0,q_80.0,q_90.0,q_100.0
f32,f32,f32,f32,f32,f32,f32,f32,f32,f32
100.069817,200.050797,300.030121,400.063232,500.00412,599.976074,699.943115,800.007996,900.036194,1000.0


In [22]:
col_2

q_10.0,q_20.0,q_30.0,q_40.0,q_50.0,q_60.0,q_70.0,q_80.0,q_90.0,q_100.0
f32,f32,f32,f32,f32,f32,f32,f32,f32,f32
100.068085,200.146133,300.243103,400.478851,500.55426,600.609314,700.744263,800.747498,900.887268,1001.0
