RUST и инжиниринг данных - часть 3
Сериализация данных
В предыдущей главе мы получали данные из OpenWeather API. Внимательные читатели заметили, что мы разобрали ответ как чистый текст, хотя ответ был в формате JSON.
Цель этого раздела - рассказать Вам о том, как получить необработанные данные и сериализовать их в структурированный формат данных, например JSON.
Мы немного углубимся в теорию чуть позже, а сейчас приступим к практике.
Сериализация
Сериализация - это процесс получения данных и их кодирования в формат, который впоследствии может быть восстановлен. Существует множество способов кодирования данных, но в основном они делятся на человекочитаемые и бинарные форматы.
CSV, JSON, XML и YAML - все это человекочитаемые форматы сериализации. Помимо них существует множество двоичных форматов, таких как Parquet, Avro и Protcol Buffers.
В конечном счете, любые данные, которые должны храниться вне памяти компьютера, требуют определенного типа сериализации.
Давайте рассмотрим сериализацию в Rust и Python.
Python
В Python с помощью модуля json мы можем сериализовать практически любую произвольную структуру данных в JSON:
In [1]:importjsonIn [2]: my_obj = [{'a': 1,'b': None},"foo","bar", ("baz","baz")]In [3]: json.dumps(my_obj)Out[3]:'[{"a": 1, "b": null}, "foo", "bar", ["baz", "baz"]]'
Вот обновленный код проекта, который сериализует ответ от API OpenWeather:
importosimportsysimportrequestsAPI_KEY = os.getenv("OWM_APPID")def get_air_pollution(lat, lon):url =f"http://api.openweathermap.org/data/2.5/air_pollution?lat={lat}&lon={lon}&appid={API_KEY}"body = requests.get(url).json()returnbodydef parse_air_pollution(body):aqi = body["list"][0]["main"]["aqi"]components = body["list"][0]["components"]return(aqi, components)if__name__ =="__main__":usage =f"Usage: python {__file__} <lat> <lon>" if notAPI_KEY:print("Please set OWM_APPID environment variable")sys.exit(1)iflen(sys.argv) !=3:print(usage)sys.exit(1)lat = sys.argv[1]lon = sys.argv[2]body = get_air_pollution(lat, lon)aqi, components = parse_air_pollution(body)print(f"Air Quality Index: {aqi}")print("Components:")fork, vincomponents.items():print(f" {k}: {v}")
Здесь необходимо отметить несколько ключевых моментов.
Во-первых, мы предполагаем, что запрос был успешным, что есть тело ответа в формате json и что оно может быть правильно разобрано. Если любое из этих предположений неверно, будет вызвано исключение, и у нас нет очевидного способа узнать, что это за исключения и какой метод может их вызвать.
def parse_air_pollution(body):aqi = body["list"][0]["main"]["aqi"]components = body["list"][0]["components"]return(aqi, components)
В процессе парсинга ответа мы «разрезаем» тело ответа на части для того, чтобы получить различные компоненты. Мы явно извлекаем ключи из словаря, предполагая, что полезная нагрузка сформирована правильно. Существуют более безопасные методы словаря, например .get(), который вернет None, если ключ отсутствует, но в нашем случае исключение оправдано, поскольку мы не можем ничего сделать с данными, если они отсутствуют.
Мы также явно не типизировали ответ от API. Это можно сделать с помощью mypy или других инструментов, таких как pydantic.
Теперь давайте посмотрим на Rust.
Rust
В Rust для reqwest нам нужно будет установить крейт serde, а также функцию json.
cargo add serde --features derivecargo add serde_jsoncargoaddreqwest --featuresjson
Поскольку Rust - типизированный язык, мы можем определить структуру, которая будет представлять ожидаемые данные. Ответ API выглядит следующим образом:
{"coord": {"lon": -122.5889,"lat": 37.9871},"list": [{"main": {"aqi": 2},"components": {"co": 168.56,"no": 0.14,"no2": 0.75,"o3": 80.11,"so2": 0.7,"pm2_5": 3.48,"pm10": 5.58,"nh3": 0},"dt": 1687308878}]}
Мы можем определить структуру, представляющую эти данные, следующим образом:
useserde::{Deserialize, Serialize};#[derive(Debug, Serialize, Deserialize)] pub struct AirPollution{pubcoord: Coord,publist:Vec<List>,}#[derive(Debug, Serialize, Deserialize)] pub struct Coord{publon:f32,publat:f32,}
#[derive(Debug, Serialize, Deserialize)] pub struct List{pubmain: Main,pubcomponents: Components,pubdt:usize,}#[derive(Debug, Serialize, Deserialize)] pub struct Main{pubaqi:u8,}#[derive(Debug, Serialize, Deserialize)] pub struct Components{pubco:f32,pubno:f32,pubno2:f32,pubo3:f32,pubso2:f32,pubpm2_5:f32,pubpm10:f32,pubnh3:f32,}
Как видите, структура полностью отражает структуру JSON. Крейт serde дает нам в данном случае большой запас гибкости, особенно в разделе «Атрибуты» и «Примеры», обязательно изучите их, они стоят того, чтобы потратить на них немного своего времени.
Крейт reqwest также подразумевает применение метода json, который автоматически десериализует тело ответа в структуру.
pub fn get_air_pollution(lat:f32, lon:f32) -> AirPollution {letapi_key = std::env::var("OWM_APPID").expect("Environment Variable OWM_APPID not set. Please set it to your OpenWeatherMap API key. https://home.openweathermap.org/api_keys",);leturl =format!("http://api.openweathermap.org/data/2.5/air_pollution?lat={}&lon={}&appid={}",lat, lon, api_key);reqwest::blocking::get(url).expect("request failed").json().expect("json failed")
Теперь наша функция возвращает не строку, а целую структуру AirPollution, а метод json в reqwest автоматически десериализует тело ответа до нужного типа.
Rust использует вывод типов для того, чтобы уменьшить количество требуемого синтаксиса. В то время как параметры функций и сигнатуры всегда требуют указания типа, локальные переменные обычно определяются компилятором.
Давайте посмотрим, как возврат типизированной структуры изменит наше взаимодействие с данными:
pub fn parse_air_pollution(body: &AirPollution) -> (&Main, &Components) {letmain = &body.list[0].main;letcomponents = &body.list[0].components;(main,components)}
Мы можем получить прямой доступ к полям, лежащим в основе структуры. В отличие от словаря Python, компилятор будет следить за тем, чтобы поля, к которым мы обращаемся, обязательно сохранялись.
Если мы добавим недостающее поле:
letfoo = &body.list[0].foo;
и выполним cargo check, мы получим ошибку:
error[E0609]: no field `foo` ontype`List`--> src/bin/ch4.rs:65:29|65 |letfoo = &body.list[0].foo;| ^^^ unknown field|= note: available fields are: `main`, `components`, `dt`For more information about this error, try `rustc --explain E0609`.error: could not compile `wxrs` (bin"ch4") due to previous error
В отличие от Python, где при попытке получить доступ к отсутствующему полю мы получим только ошибку времени выполнения (если не подключим подсказки типов с помощью mypy).
Вот полный код на языке Rust (для справки):
useserde::{Deserialize, Serialize};#[derive(Debug, Serialize, Deserialize)] pub struct AirPollution{pubcoord: Coord,publist:Vec<List>,}#[derive(Debug, Serialize, Deserialize)] pub struct Coord{publon:f32,publat:f32,}#[derive(Debug, Serialize, Deserialize)] pub struct List{pubmain: Main,pubcomponents: Components,pubdt:usize,}#[derive(Debug, Serialize, Deserialize)] pub struct Main{pubaqi:u8,}#[derive(Debug, Serialize, Deserialize)] pub struct Components{pubco:f32,pubno:f32,pubno2:f32,pubo3:f32,pubso2:f32,pubpm2_5:f32,pubpm10:f32,pubnh3:f32,}pub fn get_air_pollution(lat:f32, lon:f32) -> AirPollution {letapi_key = std::env::var("OWM_APPID").expect("Environment Variable OWM_APPID not set. Please set it to your OpenWeatherMap API key. https://home.openweathermap.org/api_keys",);leturl =format!("http://api.openweathermap.org/data/2.5/air_pollution?lat={}&lon={}&appid={}",lat, lon, api_key);reqwest::blocking::get(url).expect("request failed").json().expect("json failed")}pub fn parse_air_pollution(body: &AirPollution) -> (&Main, &Components) {letmain = &body.list[0].main;letcomponents = &body.list[0].components;(main, components)}pub fn main() {letusage =format!("Usage: {} [lat] [lon]", std::env::args().next().unwrap());letlat = std::env::args().nth(1).expect(&usage).parse::<f32>().expect(&usage);letlon = std::env::args().nth(2).expect(&usage).parse::<f32>().expect(&usage);letbody = get_air_pollution(lat, lon);let(main, components) = parse_air_pollution(&body);println!("Air Quality Index: {}", main.aqi);println!("Carbon Monoxide: {} μg/m³", components.co);println!("Nitrogen Monoxide: {} μg/m³", components.no);println!("Nitrogen Dioxide: {} μg/m³", components.no2);println!("Ozone: {} μg/m³", components.o3);println!("Sulfur Dioxide: {} μg/m³", components.so2);println!("Particulate Matter < 2.5 μm: {} μg/m³", components.pm2_5);println!("Particulate Matter < 10 μm: {} μg/m³", components.pm10);println!("Ammonia: {} μg/m³", components.nh3);}
Форматы сериализации
Стоит упомянуть, что крейт serde в Rust не содержит встроенных форматов сериализации. Вместо этого для сериализации он предоставляет специальный фреймворк. Мы установили serde_json, но Вы можете воспользоваться другими форматами, такими как serde_yaml или serde_avro.
В чем смысл?
Вы можете задаться вопросом, зачем нам определять структуру и сериализовывать тело ответа в структуру. В Python мы избегаем шаблонов, получаем прямой доступ к полям, можем немного подсказывать типы в коде, свободно использовать # type: ignore, а если наше приложение «упадет», что ж, мы просто все исправим и запустим его снова.
Вы абсолютно правы! Все это правда. Однако любой опытный программист на Python также знает обо всех способах, которыми плохо типизированный код может пойти не так.
Если Вы когда-либо создавали приложение, требующее больших вычислений и оперирующее гигабайтами данных, то наверняка сталкивались с ситуацией, когда Вам приходилось запускать приложение заново из-за сбоя. Так вот знайте, что безопасность типов не только помогает предотвратить подобные проблемы, но и повышает производительность.
Компилятор может оптимизировать код на основе известных ему типов. В Python мы можем использовать подсказки типов для того, чтобы помочь компилятору, но в конечном итоге интерпретатор Python все равно динамически разрешает типы во время выполнения. В Rust компилятор знает типы во время компиляции и может оптимизировать код уже перед выполнением запроса.
Что делает этот маленький значок &?
А, да, &. Теперь мы проникаем в самое сердце Rust. Давайте снова посмотрим на код:
pub fn parse_air_pollution(body: &AirPollution) -> (&Main, &Components) {letmain = &body.list[0].main;letcomponents = &body.list[0].components;(main, components)}
parse_air_pollution - функция, которая принимает ссылку на структуру AirPollution. Символ & - это синтаксис для создания ссылки. В Rust ссылки - это способ передачи значения в функцию без передачи права собственности на него. Это ключевое понятие в Rust, именно оно позволяет Rust гарантировать безопасность памяти.
В Python значения передаются с помощью счетчиков. Каждый раз, когда Вы используете переменную, сборщик мусора Python отслеживает, сколько раз она была использована. Всякий раз, когда существует функция, использовавшая ссылку, счетчик уменьшается. Время от времени сборщик мусора запускается и очищает все неиспользуемые ссылки.
В Rust сборщик мусора отсутствует. Вместо этого компилятор отслеживает время жизни каждой переменной. Когда переменная выходит из области видимости, компилятор автоматически освобождает память, связанную с ней.
Это означает, что Вы не сможете использовать переменную после передачи прав собственности. Для более глубокого погружения в концепцию владения прочитайте Книгу Rust.
Например, если бы мы попытались вывести значение body после присвоения, компилятор выдал бы cледующую ошибку:
fn parse_air(body: AirPollution) {letfoo = body;println!("{:?}", body);}error[E0382]: borrow of moved value: `body`--> src/bin/ch4.rs:71:22|69 | fn parse_air(body: AirPollution) {| ---- move occurs because `body` has type `AirPollution`, which does not implement the `Copy` trait70 | let foo = body;| ---- value moved here71 | println!("{:?}", body);| ^^^^ value borrowed here after move
Объяснение всех тонкостей владения и ссылок выходит за рамки этой статьи, но в любом случае очень важно понимать, что компилятор Rust отслеживает время жизни каждой переменной и не позволит Вам использовать переменную после того, как она была перемещена.
Но мы можем использовать ссылку на переменную. Это позволяет сохранить базовые данные в том же месте в памяти, но передать их в функцию как ссылку на исходное значение.
fn parse_air(body: &AirPollution) {letfoo = body;println!("{:?}", body);}
Это дает некоторые преимущества, когда речь идет об обработке больших объемов данных.
В Python не всегда понятно, когда данные копируются, перемещаются или на них ссылаются. В Rust копирование кода происходит явно. Если бы мы не хотели заимствовать ссылку в приведенном выше коде, мы могли бы его просто скопировать.
fn parse_air(body: AirPollution) {letfoo = body.clone();println!("{:?}", body);}
Чтобы приведенный выше код сработал, нам также необходимо реализовать функцию Clone для структуры AirPollution и всех ее полей:
#[derive(Debug, Clone, Deserialize)] pub struct AirPollution{...
Понимание владения, ссылок и заимствований может оказаться нелегким испытанием для начинающих программистов на Rust, привыкших к динамически-типизированным языкам, но со временем все обязательно получится.
Производительность
Чтобы проверить наш код, мы изменим его таким образом, чтобы получать прогноз не на один день, а на гораздо больший срок, что увеличит полезную нагрузку с 0,5 кб до примерно 13 кб.
В Python мы изменим url, а затем выполним итерацию по каждому элементу в полученном списке.
def get_air_pollution(lat, lon):url =f"http://api.openweathermap.org/data/2.5/air_pollution/forecast?lat={lat}&lon={lon}&appid={API_KEY}"body = requests.get(url).json()returnbodydef parse_air_pollution(body):res = []print(body)forrowinbody["list"]:res.append((row["main"]["aqi"], row["components"], row["dt"]))returnresdef print_air_pollution(main, components, dt):print("---")print(f"Air pollution forecast for {dt}")print(f"Air quality index: {main}")print("Components:")fork, vincomponents.items():print(f" {k}: {v}")
В Rust мы также изменим url и используем общий шаблон iter().map().collect()..
pub fn get_air_pollution(lat:f32, lon:f32) -> AirPollution {letapi_key = std::env::var("OWM_APPID").expect("Environment Variable OWM_APPID not set. Please set it to your OpenWeatherMap API key. https://home.openweathermap.org/api_keys",);leturl =format!("http://api.openweathermap.org/data/2.5/air_pollution/forecast?lat={}&lon={}&appid={}",lat, lon, api_key);reqwest::blocking::get(url).expect("request failed").json().expect("json failed")}pub fn parse_air_pollution(body: AirPollution) ->Vec<(Main, Components,usize)> {body.list.iter().map(|x| (x.main, x.components, x.dt)).collect()}
Вот результаты сравнений:
И снова мы видим улучшение производительности в 1,7 раза (примерно на 58%).
Offline сравнение
Сравнение при сетевом подключении может быть немного сомнительным. Это также затрудняет тестирование все больших и больших полезных нагрузок, поэтому мы создадим большой файл полезной нагрузки и используем его для оффлайн бенчмарка.
Я создал JSON-файл размером 9 мб, который поддерживает полезную нагрузку из API OpenWeather, а также автономные версии кода на Rust и Python для чтения из локального файла. Код для обеих версий можно найти в репозитории примеров:
under wxpy/wxpy/ch4/serialized_offline_benchmark.py and wxrs/src/bin/ch4_offline_benchmark.rs.
Результаты офлайн-сравнения:
|
Команда |
Среднее значение [мс] |
Min [мс |
Max [мс] |
Относительное значение |
|---|---|---|---|---|
|
|
47.0 ± 0.6 |
46.5 |
50.6 |
1.00 |
|
|
248.0 ± 22.4 |
210.8 |
278.2 |
5.27 ± 0.48 |
При таких больших нагрузках Rust работает в два раза быстрее, чем Python.
Преобразование данных с помощью Polars
В этой главе мы рассмотрим, как преобразовывать данные с помощью Polars в Python и Rust.
Polars - это «молниеносно быстрая библиотека DataFrame», доступная как в Python, так и в Rust. Она похожа на pandas, но обладает меньшими возможностями, однако поддерживает широкий спектр распространенных задач преобразования.
Она в несколько раз быстрее, чем pandas, и является оптимальным выбором для решения задач преобразования данных.
Рекомендую начать изучение этой библиотеки с официальной документации.
Рассмотрим ключевые различия между синтаксисом на Python и на Rust.
Python
import os
import polars as pl
script_path = os.path.dirname(os.path.realpath(__file__))
bird_path = os.path.join(script_path, "../../../lib/PFW_2016_2020_public.csv")
codes_path = os.path.join(script_path, "../../../lib/species_code.csv")
columns = [
"LATITUDE",
"LONGITUDE",
"SUBNATIONAL1_CODE",
"Month",
"Day",
"Year",
"SPECIES_CODE",
"HOW_MANY",
"VALID",
]
birds = pl.read_csv(
bird_path,
columns=columns,
new_columns=[s.lower() for s in columns],
)
codes = pl.read_csv(codes_path).select(
[
pl.col("SPECIES_CODE").alias("species_code"),
pl.col("PRIMARY_COM_NAME").alias("species_name"),
]
)
birds_df = (
birds.select(
pl.col(
[
"latitude",
"longitude",
"subnational1_code",
"month",
"day",
"year",
"species_code",
"how_many",
"valid",
]
)
)
.filter(pl.col("valid") == 1)
.groupby(["subnational1_code", "species_code"])
.agg(
[
pl.sum("how_many").alias("total_species"),
pl.count("how_many").alias("total_sightings"),
]
)
.sort("total_species", descending=True)
.join(codes, on="species_code", how="inner")
)
print(birds_df)
Код на Python лаконичен, столбцы можно выбирать в виде списка строк, функция сортировки принимает простой аргумент по убыванию, а общий API очень прост.
Я попытался реализовать ту же логику в pandas. Несмотря на значительную схожесть, есть несколько отличий, например, в том, как мы фильтруем достоверные результаты:
import os
import pandas as pd
script_path = os.path.dirname(os.path.realpath(__file__))
bird_path = os.path.join(script_path, "../../../lib/PFW_2016_2020_public.csv")
codes_path = os.path.join(script_path, "../../../lib/species_code.csv")
# adding usecols reducing memory usage and runtime from 13s to 7s
birds = pd.read_csv(
bird_path,
usecols=[
"LATITUDE",
"LONGITUDE",
"SUBNATIONAL1_CODE",
"Month",
"Day",
"Year",
"SPECIES_CODE",
"HOW_MANY",
"VALID",
],
).rename(columns=lambda x: x.lower())
codes = pd.read_csv(codes_path)[["SPECIES_CODE", "PRIMARY_COM_NAME"]].rename(
columns={"SPECIES_CODE": "species_code", "PRIMARY_COM_NAME": "species_name"}
)
birds = birds[
[
"latitude",
"longitude",
"subnational1_code",
"month",
"day",
"year",
"species_code",
"how_many",
"valid",
]
]
birds = birds[birds["valid"] == 1]
birds = (
birds.groupby(["subnational1_code", "species_code"])
.agg(total_species=("how_many", "sum"), total_sightings=("how_many", "count"))
.reset_index()
.sort_values("total_species", ascending=False)
)
birds = pd.merge(birds, codes, on="species_code", how="inner")
print(birds)
Теперь сравним все вышеперечисленное с кодом на Rust.
Rust
use polars::prelude::*;
use std::env;
fn main() {
let current_dir = env::current_dir().expect("Failed to get current directory");
let bird_path = current_dir.join("../lib/PFW_2016_2020_public.csv");
let codes_path = current_dir.join("../lib/species_code.csv");
let cols = vec![
"LATITUDE".into(),
"LONGITUDE".into(),
"SUBNATIONAL1_CODE".into(),
"Month".into(),
"Day".into(),
"Year".into(),
"SPECIES_CODE".into(),
"HOW_MANY".into(),
"VALID".into(),
];
let birds_df = CsvReader::from_path(bird_path)
.expect("Failed to read CSV file")
.has_header(true)
.with_columns(Some(cols.clone()))
.finish()
.unwrap()
.lazy();
let mut codes_df = CsvReader::from_path(codes_path)
.expect("Failed to read CSV file")
.infer_schema(None)
.has_header(true)
.finish()
.unwrap();
codes_df = codes_df
.clone()
.lazy()
.select([
col("SPECIES_CODE").alias("species_code"),
col("PRIMARY_COM_NAME").alias("species_name"),
])
.collect()
.unwrap();
let birds_df = birds_df
.rename(cols.clone(), cols.into_iter().map(|x| x.to_lowercase()))
.filter(col("valid").eq(lit(1)))
.groupby(["subnational1_code", "species_code"])
.agg(&[
col("how_many").sum().alias("total_species"),
col("how_many").count().alias("total_sightings"),
])
.sort(
"total_species",
SortOptions {
descending: true,
nulls_last: false,
multithreaded: true,
},
)
.collect()
.unwrap();
let joined = birds_df
.join(
&codes_df,
["species_code"],
["species_code"],
JoinType::Inner,
None,
)
.unwrap();
println!("{}", joined);
}
В Rust код на 75% длиннее, а синтаксис более многословен. Здесь много вызовов unwrap для обработки ошибок, хотя в реальном приложении некоторые из них можно заменить на ?
Функция sort принимает структуру SortOptions, которая немного более многословна. В целом, API очень похож.
Сравнение
Давайте рассмотрим несколько сравнение polars на Python и Rust, а также аналогичный код на Pandas.
|
Команда |
Среднее значение[мс] |
Min [ms] |
Max [ms] |
Относительное значение |
|---|---|---|---|---|
|
../wxrs/target/release/ch5 |
473.9 ± 5.8 |
461.3 |
480.7 |
1.00 |
|
python ../wxpy/wxpy/ch5/ch5.py |
764.8 ± 26.8 |
732.9 |
815.2 |
1.61 ± 0.06 |
|
python ../wxpy/wxpy/ch5/ch5_pandas.py |
5644.0 ± 39.2 |
5584.9 |
5710.0 |
11.91 ± 0.17 |
Версия на Rust снова самая быстрая. Код на Python-polars в 1,6 раза медленнее, чем код на Rust, код на Pandas самый медленный: на его выполнение уходит более 5 секунд, в то время как обе версии на Polars занимают менее 1 секунды.




