Использование Trino совместно с Dataproc
Trino (ранее Presto) - это распределенный механизм SQL-запросов, предназначенный для обработки больших массивов данных, распределенных по одному или нескольким разнородным источникам данных. Trino может выполнять запросы к Hive, MySQL, Kafka и другим источникам данных через коннекторы. В этом руководстве показано, как:
- установить Trino на кластер Dataproc
- Запрашивать публичные данные с помощью клиента Trino, установленного на локальной машине, который взаимодействует со службой Trino
- Выполнять запросы из Java-приложения, которое взаимодействует с сервисом Trino на через драйвер Trino Java JDBC.
Цели
- создание кластера Dataproc с установленным Trino
- подготовка данных. В этом руководстве использованы данные о поездах на такси в Чикаго, доступные в BigQuery.
- Извлечь данные из BigQuery
- Загрузить данные в Cloud Storagе в формате файлов CSV
- Преобразовать данные:
- Представьте данные в виде внешней таблицы Hive, чтобы сделать их доступными для запросов в Trino
- Преобразуйте данные из формата CSV в формат Parquet (для ускорения обработки запросов)
- Отправьте запросы Trino CLI или кода приложения с помощью туннеля SSH или драйвера Trino JDBC координатору Trino, работающему на кластере
- Проверьте журналы Trino через веб-интерфейс Trino.
Затраты
В этом документе используются следующие платные компоненты Google Cloud:
Для расчета примерной стоимости проекта используйте специальный калькулятор.
Перед началом работы
Если Вы еще не сделали этого, создайте проект Google Cloud и бакет Cloud Storage (для хранения данных, используемых в этом руководстве)
Настройка проекта
- В консоли Google Cloud на странице выбора проекта выберите или создайте проект Google Cloud.
NB: Если Вы не планируете хранить ресурсы, созданные в этой процедуре, создайте отдельный проект, не выбирайте существующий. После выполнения этих шагов Вы сможете удалить проект, удалив все связанные с ним ресурсы.
- Выберите Отборщик проектов (Project selector)
- Убедитесь в том, что для Вашего проекта Cloud включен биллинг.
- Включите API Dataproc, Compute Engine, Cloud Storage и BigQuery.
- Включите API
- Установите the Google Cloud CLI. Для инициализации gcloud CLI выполните следующую команду:
gcloud init
Создание бакета для работы с данными, используемыми в данном руководстве.
- В консоли Google Cloud перейдите на страницу Cloud Storage Buckets.
- Перейдите на страницу бакета
- Щелкните на Создать бакет.
- На открывшейся странице введите название бакета. Нажмите Продолжить.
- Присвоить бакету имя.
-
Выберите, где Вы будете хранить данные:
- Выберите опцию Тип локации.
- Выберите опцию Локация.
- На вкладке Выбрать класс хранения для Ваших данных выберите класс хранения.
- На вкладке Выбрать способ контроля доступа к объектам данных выберите опцию Контроль доступа.
- В качестве Расширенных настроек (по желанию) укажите метод шифрования, политику хранения или метки бакета.
- Нажмите Создать.
Создание кластера Dataproc
Создайте кластер, выполнив команды, указанные ниже.
Создайте кластер Dataproc, используя флаг optional-components (доступен на образе версии 2.1 и новее) для установки дополнительного компонента Trino на кластер и флаг enable-component-gateway для включения шлюза компонентов, который позволит получить доступ к веб-интерфейсу Trino из консоли Google Cloud.
-
Установите переменные окружения:
- PROJECT: ID проекта
- BUCKET_NAME: имя бакета облачного хранилища, созданного на этапе Перед началом работы
- REGION: регион, в котором будет создан кластер, например, "us-west1"
- WORKERS: рекомендуем использовать 3 - 5 рабочих узла
export PROJECT=project-id export WORKERS=number export REGION=region export BUCKET_NAME=bucket-name
- Для создания кластера запустите Google Cloud CLI.
gcloud beta dataproc clusters create trino-cluster \
--project=${PROJECT} \
--region=${REGION} \
--num-workers=${WORKERS} \
--scopes=cloud-platform \
--optional-components=TRINO \
--image-version=2.1 \
--enable-component-gateway
Мы не рекомендуем создавать высокодоступный кластер (кластер с несколькими мастер - узлами), поскольку координатор запускается на мастер - узле 0, все остальные мастер - узлы будут простаивать.
Подготовка данных
-
Экспортируйте
набор данных опоездках на такси в Чикаго в облачное хранилище в формате CSV, а затем создайте внешнюю таблицу Hive для того, чтобы на эти данные можно было ссылаться. - На локальной машине выполните команду, которая позволит импортировать данные о такси из BigQuery в виде CSV-файлов без заголовков бакет облачного хранилища, которое Вы создали на этапе Перед началом работы.
bq --location=us extract --destination_format=CSV \
--field_delimiter=',' --print_header=false \
"bigquery-public-data:chicago_taxi_trips.taxi_trips" \
gs://${BUCKET_NAME}/chicago_taxi_trips/csv/shard-*.csv3.
Создайте внешнюю таблицу Hive под названием chicago_taxi_trips_csv.
gcloud dataproc jobs submit hive \
--cluster trino-cluster \
--region=${REGION} \
--execute "
CREATE EXTERNAL TABLE chicago_taxi_trips_csv(
unique_key STRING,
taxi_id STRING,
trip_start_timestamp TIMESTAMP,
trip_end_timestamp TIMESTAMP,
trip_seconds INT,
trip_miles FLOAT,
pickup_census_tract INT,
dropoff_census_tract INT,
pickup_community_area INT,
dropoff_community_area INT,
fare FLOAT,
tips FLOAT,
tolls FLOAT,
extras FLOAT,
trip_total FLOAT,
payment_type STRING,
company STRING,
pickup_latitude FLOAT,
pickup_longitude FLOAT,
pickup_location STRING,
dropoff_latitude FLOAT,
dropoff_longitude FLOAT,
dropoff_location STRING)
ROW FORMAT DELIMITED
FIELDS TERMINATED BY ','
STORED AS TEXTFILE
location 'gs://${BUCKET_NAME}/chicago_taxi_trips/csv/';"
- Проверьте создание внешней таблицы Hive. Если Ваш кластер состоит из 3 или менее узлов, этот запрос может занять несколько минут. Вы можете пропустить этот шаг и проверить только создание конечного файла Parquet, описанного ниже.
gcloud dataproc jobs submit hive \
--cluster trino-cluster \
--region=${REGION} \
--execute "SELECT COUNT(*) FROM chicago_taxi_trips_csv;"6.
- Загрузите данные из таблицы Hive CSV в таблицу Hive Parquet.
gcloud dataproc jobs submit hive \
--cluster trino-cluster \
--region=${REGION} \
--execute "
INSERT OVERWRITE TABLE chicago_taxi_trips_parquet
SELECT * FROM chicago_taxi_trips_csv;"
- Убедитесь в том, что данные загрузились правильно.
gcloud dataproc jobs submit hive \
--cluster trino-cluster \
--region=${REGION} \
--execute "SELECT COUNT(*) FROM chicago_taxi_trips_parquet;"
Выполните запросы
Вы можете запускать запросы локально из Trino CLI или из приложения.
Запросы Trino CLI
В этом разделе показано, как выполнять запросы к набору данных Hive Parquet taxi с помощью Trino CLI.
- Выполните следующую команду, которая поможет Вам подключиться по SSH к главному узлу кластера. Во время выполнения команды локальный терминал перестанет отвечать на запросы.
gcloud compute ssh trino-cluster-m
- В окне терминала SSH на главном узле кластера запустите Trino CLI, который подключается к серверу Trino, запущенному на главном узле.
trino --catalog hive --schema default
-
В
trino:defaultубедитесь в том, что Trino может найти таблицы Hive.
show tables; Table ‐‐‐‐‐‐‐‐‐‐‐‐‐‐‐‐‐‐‐‐‐‐‐‐‐‐‐‐‐ chicago_taxi_trips_csv chicago_taxi_trips_parquet (2 rows)
-
Выполните запросы из промта
trino:defaultи сравните производительность запросов к данным Parquet и CSV.
Данные Parquet
select count(*) from chicago_taxi_trips_parquet where trip_miles > 50; _col0 ‐‐‐‐‐‐‐‐ 117957 (1 row) Query 20180928_171735_00006_2sz8c, FINISHED, 3 nodes Splits: 308 total, 308 done (100.00%) 0:16 [113M rows, 297MB] [6.91M rows/s, 18.2MB/s]
Данные CSV
select count(*) from chicago_taxi_trips_csv where trip_miles > 50; _col0 ‐‐‐‐‐‐‐‐ 117957 (1 row) Query 20180928_171936_00009_2sz8c, FINISHED, 3 nodes Splits: 881 total, 881 done (100.00%) 0:47 [113M rows, 41.5GB] [2.42M rows/s, 911MB/s]
Приложение Java
Запросы из приложения Java с помощью драйвера Trino Java JDBC: 1. Скачайте драйвер Trino Java JDBC 1. В Maven pom.xml добавьте зависимость trino-jdbc.
<dependency> <groupId>io.trino</groupId> <artifactId>trino-jdbc</artifactId> <version>376</version> </dependency>
Образец кода Java
package dataproc.codelab.trino;
import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.sql.Statement;
import java.util.Properties;
public class TrinoQuery {
private static final String URL = "jdbc:trino://trino-cluster-m:8080/hive/default";
private static final String SOCKS_PROXY = "localhost:1080";
private static final String USER = "user";
private static final String QUERY =
"select count(*) as count from chicago_taxi_trips_parquet where trip_miles > 50";
public static void main(String[] args) {
try {
Properties properties = new Properties();
properties.setProperty("user", USER);
properties.setProperty("socksProxy", SOCKS_PROXY);
Connection connection = DriverManager.getConnection(URL, properties);
try (Statement stmt = connection.createStatement()) {
ResultSet rs = stmt.executeQuery(QUERY);
while (rs.next()) {
int count = rs.getInt("count");
System.out.println("The number of long trips: " + count);
}
}
} catch (SQLException e) {
e.printStackTrace();
}
}
}Ведение журнала и мониторинг
Ведение журнала
Журналы Trino находятся по адресу /var/log/trino/ на главном и рабочих узлах кластера.
Мониторинг
Trino раскрывает информацию о времени выполнения запросов через специальные таблицы. Для просмотра данных таблицы времени выполнения выполните следующий запрос в сеансе Trino (из промта trino:default):
select * FROM system.runtime.nodes;
Очистка
После завершения обучения Вы можете очистить созданные ресурсы. В следующих разделах описаны способы удаления и отключения этих ресурсов.
Удаление проекта
Самый простой способ удалить биллинг - удалить проект, который Вы создали в учебных целях.
Внимание: последствия удаления проекта:
- Все, что было в проекте, удаляется. Если Вы использовали существующий проект для выполнения задач, описанных в этом документе, при его удалении удаляются и все другие работы, выполненные в рамках проекта.
- Пользовательские идентификаторы проекта теряются. При создании этого проекта Вы создали пользовательский идентификатор проекта, который, возможно, захотите использовать в будущем. Чтобы сохранить URL-адреса, использующие идентификатор проекта, например URL-адрес appspot.com, удалите выбранные ресурсы внутри проекта, не удаляйте весь проект.
Если Вы планируете изучить другие быстрые запуски, повторное использование проекта поможет Вам избежать превышения лимита проектной квоты.
- В консоли Google Cloud перейдите на страницу Управление ресурсами.
- Управление ресурсами
- В списке проектов выберите проект, который вы хотите удалить, а затем нажмите кнопку Удалить.
- В диалоговом окне введите идентификатор проекта, а затем нажмите Закрыть, чтобы удалить проект.
Удаление кластера
gcloud dataproc clusters delete --project=${PROJECT} trino-cluster \
--region=${REGION}Удаление бакета
gcloud storage rm gs://${BUCKET_NAME} --recursive





