-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathpython.qmd
More file actions
167 lines (128 loc) · 4.62 KB
/
Copy pathpython.qmd
File metadata and controls
167 lines (128 loc) · 4.62 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
---
title: "Démonstration Apache Arrow en Python"
format: html
execute:
cache: true
---
```{python}
#| tags: [parameters]
#| echo: false
chemin_dataset = "C:\\Users\\afeldmann\\Desktop\\openamir"
```
# Introduction
Dans cette démonstration, nous allons explorer comment Apache Arrow peut être utilisé pour traiter des fichiers CSV et Parquet, même lorsque ces derniers sont trop volumineux pour tenir en mémoire. Nous allons montrer :
1. Comment charger un fichier CSV en mémoire avec Arrow et effectuer un tri.
2. Comment utiliser l'API `pyarrow.dataset` pour travailler avec un fichier Parquet trop grand pour la mémoire, appliquer des filtres et des agrégations, et convertir les résultats en un DataFrame Pandas.
```{python}
#| eval: false
! pip install pyarrow
! pip install pandas
```
# Chargement d'un fichier CSV avec Arrow
```{python}
import pyarrow.csv as csv
csv_file = "data1.csv"
table = csv.read_csv(csv_file)
# Trier les données par une colonne spécifique
sorted_table = table.sort_by("name")
print(sorted_table)
```
# Traitement d'un fichier Parquet trop grand pour la mémoire
```{python}
import pyarrow.dataset as ds
import pyarrow.compute as pc
import pandas as pd
```
```{python}
#| echo: false
#| results: false
emplacement_dataset = chemin_dataset
```
```{python}
#| echo: true
#| eval: false
emplacement_dataset = "[emplacement du dataset]"
```
```{python}
dataset = ds.dataset(emplacement_dataset, format="parquet", partitioning="hive")
# Appliquer des filtres et des agrégations
filtered = dataset.filter(pc.field("BEN_SEX_COD") == "1")
```
## La lazy evaluation :
Le code de la cellule précédente était instantanée et celui qui vient produit un résultat, il va donc le déclencher
```{python}
filtered.count_rows()
filtered.head(20)
```
## Faire une aggregation sur une table dépassant la mémoire avec Acero :
```{python}
import pyarrow.dataset as ds
import pyarrow.acero as acero
import pyarrow.compute as pc
dataset = ds.dataset(
emplacement_dataset,
format="parquet",
partitioning="hive"
)
plan = acero.Declaration(
"scan",
acero.ScanNodeOptions(dataset, columns=["FLT_PAI_MNT", "annee", "mois"])
)
plan = acero.Declaration(
"project",
acero.ProjectNodeOptions(
expressions=[
pc.field("FLT_PAI_MNT"),
pc.field("annee"),
pc.field("mois"),
],
names=["FLT_PAI_MNT", "annee", "mois"]
),
inputs=[plan]
)
plan = acero.Declaration(
"aggregate",
acero.AggregateNodeOptions(
aggregates=[
("FLT_PAI_MNT", "hash_sum", pc.ScalarAggregateOptions(skip_nulls=False, min_count=0), "FLT_PAI_MNT")
],
keys=["annee", "mois"]
),
inputs=[plan]
)
```
```{python}
plan
```
Pour l'instant, c'est une requête, qui fait exactement la même requête que celle que nous avions sous R auparavant.
```{python}
table = plan.to_table()
print(table)
```
On peut l'exécuter à l'aide de la fonction `to_table()`.
## Conversion en pandas
```{python}
# Convertir en DataFrame Pandas
pandas_df = table.to_pandas()
print(pandas_df)
```
## Version polars
```{python}
import polars as pl
plan = (
pl.scan_parquet(emplacement_dataset)
.select(["FLT_PAI_MNT", "annee", "mois"])
.group_by(["annee", "mois"])
.agg(pl.col("FLT_PAI_MNT").sum())
)
print(plan)
```
L'on peut également utiliser polars en backend, qui est un moteur alternatif, basé sur le format de donnée `arrow`, mais qui utilise un moteur alternatif à Acero pour exécuter ses requêtes.
Comme Acero, Polars attend qu'on lui ordonne d'exécuter ses requêtes.
Comme sous Acero, on peut éviter de tout charger en mémoire d'un seul coup sous Polars. Il faut pour cela préciser qu'on utilise l'engine "streaming".
```{python}
plan.collect(engine="streaming")
```
Dès lors que l'on atteint la dizaine de secondes pour des requêtes d'études statistiques, l'on peut se dire que l'objectif est rempli et que le benchmark comparatif n'importe pas beaucoup pour des chercheurs. Si l'on reste dans le cadre de requêtes accomplies par l'homme, l'important est la **lisibilité du frontend** ainsi que la **stabilité du framework**. De ce point de vue, l'implémentation R d'Acero que l'on a précédemment vue, est très satisfaisante (et a l'avantage d'être maintenue par Apache).
Il est clair que, contrairement à l'implémentation R d'Acero, la lisibilité d'`ExecPlan` sous `pyarrow` est assez moyenne (ça reste quelque chose d'assez bas niveau).
Dans ce contexte, polars est sans doute un bon apport à la data-analyse en Python. Elle apporte l'équivalent de `open_dataset()` en Python, à savoir une interface sous la forme de pipelines.