Соединение и счетчик Hadoop MapReduce на примере

⚡ Умное резюме

Объединение данных в MapReduce объединяет два больших набора данных по общему ключу либо внутри маппера, либо внутри редуктора, а счетчики MapReduce собирают статистику о задании, что позволяет измерять, а не угадывать наличие некорректных записей.

  • 🔘 Основные моменты подключения: Меньший из двух наборов данных распределяется по каждому узлу данных и используется в качестве поискового.
  • ☑️ Соединение на стороне карты: Для корректной работы функции map необходимо, чтобы каждый входной поток был разбит на части, равномерно распределен и отсортирован по ключу объединения.
  • Соединение со стороной Reduce: Разделение на разделы не требуется, поскольку каждый кортеж, имеющий общий ключ соединения, попадает в один и тот же редуктор.
  • 🧪 Реализованный пример: Файлы DeptName.txt и DeptStrength.txt копируются в HDFS и объединяются по идентификатору Dept_ID с помощью упакованного JAR-файла.
  • 🇧🇷 Типы счетчиков: В каждом задании предусмотрено пять встроенных групп счетчиков, а определяемые пользователем счетчики объявляются как Java перечисление.
  • ⚠️ Противодействие: Увеличение счетчика при каждой отсутствующей или некорректной записи превращает проблемы с качеством данных в числовое значение в отчете о задании.

Учебное пособие по Hadoop MapReduce с использованием операций объединения и подсчета с подробным разбором на примере.

Что такое Join в MapReduce?

Операция объединения MapReduce используется для объединения двух больших наборов данных. Однако этот процесс требует написания большого количества кода для выполнения самой операции объединения. Объединение двух наборов данных начинается со сравнения размеров каждого набора данных. Если один набор данных меньше другого, то меньший набор данных распределяется по всем узлам данных в кластере.

После присоединения Уменьшение карты В распределенной системе либо Mapper, либо Reducer используют меньший набор данных для поиска соответствующих записей в большом наборе данных, а затем объединяют эти записи для формирования выходных данных.

Типы присоединения

В зависимости от места фактического выполнения объединения, объединения в Hadoop классифицируются на два типа.

  1. соединение со стороны карты — Когда объединение выполняется маппером, это называется объединением на стороне маппера. В этом случае объединение выполняется до того, как данные фактически будут обработаны функцией маппера. Обязательно, чтобы входные данные для каждого маппера имели вид раздела и были отсортированы. Кроме того, должно быть одинаковое количество разделов, и они должны быть отсортированы по ключу объединения.
  2. Соединение с уменьшенной стороной — Когда объединение выполняется редуктором, это называется объединением на стороне редуктора. В этом случае нет необходимости иметь набор данных в структурированном виде (или секционированном). Здесь обработка на стороне карты выдает ключ объединения и соответствующие кортежи обеих таблиц. В результате этой обработки все кортежи с одинаковым ключом объединения попадают в один и тот же редуктор, который затем объединяет записи с одинаковым ключом объединения.

Общий процесс соединений в Hadoop показан на диаграмме ниже.

Схема процесса сравнения объединения на стороне карты и объединения на стороне редукции в Hadoop.
Типы объединений в Hadoop MapReduce

После того, как два варианта стали понятны, в следующем разделе рассматривается операция объединения с уменьшением размера (reduce-side join) для двух небольших файлов отдела.

Как объединить два набора данных: пример MapReduce

Имеется два набора данных в двух разных файлах (показаны ниже). Ключ Dept_ID является общим для обоих файлов. Цель состоит в том, чтобы использовать MapReduce Join для объединения этих файлов.

Первый входной файл содержит идентификаторы отделов наряду с названиями отделов.

Файл 1
Второй входной файл содержит идентификаторы отделов, а также значения численности сотрудников в каждом отделе.

Файл 2

Входной сигнал: Входные данные представляют собой текстовые файлы: DeptName.txt и DeptStrength.txt.

Загрузите исходные файлы отсюда

Убедитесь, что у вас есть Hadoop Установлено. Прежде чем приступить к фактическому процессу примера объединения MapReduce, смените пользователя на 'hduser' (идентификатор, использованный при настройке Hadoop; вы можете переключиться на идентификатор пользователя, использованный при настройке Hadoop).

su - hduser_

В командной строке отображается имя учетной записи Hadoop, как показано ниже.

Терминал после переключения на учетную запись hduser с помощью команды su

Шаг 1) Скопируйте zip-файл в выбранное вами место.

Загруженный архив MapReduceJoin помещен в выбранную рабочую директорию.

Шаг 2) Распакуйте ZIP-файл

sudo tar -xvf MapReduceJoin.tar.gz

ЭксtracНазвания файлов прокручиваются по мере того, как tar распаковывает архив.

Консольный список файлов, напримерtracted из MapReduceJoin.tar.gz

Шаг 3) Перейдите в каталог MapReduceJoin/.

cd MapReduceJoin/

После перехода в каталог MapReduceJoin появляется приглашение командной строки.

Шаг 4) Запустить Hadoop

$HADOOP_HOME/sbin/start-dfs.sh
$HADOOP_HOME/sbin/start-yarn.sh

Оба скрипта выводят на экран имена запускаемых ими демонов.

Сообщения при запуске от скриптов демонов HDFS и YARN.

Шаг 5) DeptStrength.txt и DeptName.txt — это входные файлы, используемые в этом примере программы MapReduce Join.

Эти файлы необходимо скопировать в HDFS используя приведенную ниже команду:

$HADOOP_HOME/bin/hdfs dfs -copyFromLocal DeptStrength.txt DeptName.txt /

Оба входных текстовых файла скопированы в корневой каталог HDFS.

Шаг 6) Запустите программу, используя команду ниже:

$HADOOP_HOME/bin/hadoop jar MapReduceJoin.jar MapReduceJoin/JoinDriver/DeptStrength.txt /DeptName.txt /output_mapreducejoin

Сначала команда выводится на экран, после чего задание сообщает о ходе выполнения на консоли.

Запуск JAR-файла MapReduceJoin из командной строки

Консольный вывод tracотслеживание хода выполнения задания объединения MapReduce

Шаг 7) После выполнения выходной файл (с именем 'part-00000') будет сохранен в каталоге /output_mapreducejoin в HDFS.

Результаты можно увидеть с помощью интерфейса командной строки.

$HADOOP_HOME/bin/hdfs dfs -cat /output_mapreducejoin/part-00000

Присоединённые к отделу записи распечатаны из HDFS командой cat.

Результаты также можно увидеть через веб-интерфейс, как

Главная страница веб-интерфейса Hadoop используется для доступа к файловому браузеру.

Теперь выберите «Просмотреть файловую систему» ​​и перейдите к адресу /output_mapreducejoin

Просмотр представления файловой системы HDFS в каталоге output_mapreducejoin

Открыть часть-r-00000

Выбор выходного файла part-r-00000 в окне браузера.

Результаты показаны

В браузере отображаются строки с названием объединенного отдела и численностью сотрудников отдела.

ПРИМЕЧАНИЕ: Обратите внимание, что перед следующим запуском этой программы вам необходимо будет удалить выходной каталог /output_mapreducejoin.

$HADOOP_HOME/bin/hdfs dfs -rm -r /output_mapreducejoin

Альтернативой является использование другого имени для выходного каталога.

Объединения (Join) показывают, как выглядят данные. Счетчики (Counters), которые будут рассмотрены далее, показывают, как вела себя задача, создавшая эти данные.

Что такое счетчик в MapReduce?

Счетчик в MapReduce — это механизм, используемый для сбора и измерения статистической информации о заданиях и событиях MapReduce. Счетчики хранят информацию о... tracВ MapReduce счетчики собирают различную статистику заданий, такую ​​как количество выполненных операций и ход их выполнения. Счетчики используются для диагностики проблем в MapReduce.

Счетчики Hadoop аналогичны помещению сообщения журнала в код карты или сокращения. Эта информация может быть полезна для диагностики проблем при обработке заданий MapReduce.

Как правило, в Hadoop эти счетчики определяются в программе (map или reduce) и увеличиваются во время выполнения при возникновении определенного события или условия (специфичного для этого счетчика). Очень хорошее применение счетчиков Hadoop — это track допустимых и недопустимых записей из входного набора данных.

Типы счетчиков MapReduce

В основном существует 2 типа счетчиков MapReduce.

  1. Встроенные счетчики Hadoop: Есть несколько встроенных счетчиков Hadoop, которые существуют для каждого задания. Ниже представлены встроенные группы счетчиков-
    • Счетчики задач MapReduce — Собирает информацию, специфичную для задачи (например, количество входных записей), во время ее выполнения.
    • Счетчики файловой системы — Собирает информацию, такую ​​как количество байтов, прочитанных или записанных задачей.
    • Счетчики FileInputFormat — Собирает информацию о количестве байтов, прочитанных через FileInputFormat.
    • Счетчики FileOutputFormat — Собирает информацию о количестве байтов, записанных с помощью FileOutputFormat.
    • Счетчики вакансий — Эти счетчики регистрируют статистические данные по всей работе, например, количество задач, запущенных в рамках данной работы.
  2. Счетчики, определяемые пользователем: Помимо встроенных счетчиков, пользователь может определять собственные счетчики, используя аналогичные функции, предоставляемые языками программирования. Например, в JavaПеречисление (enum) используется для определения пользовательских счетчиков.

💡 Примечание к версии: Счетчики вакансий обслуживались отделом вакансий.TracВ MRv1 эта роль принадлежит MapReduce ApplicationMaster, поэтому имена счетчиков сохраняются, но компонент, который их сообщает, изменился.

В задании нельзя установить неограниченное количество счетчиков. mapreduce.job.counters.max По умолчанию этот параметр ограничивает общее количество заданий на одно задание 120, и задание, в котором указано больше, завершается с ошибкой. LimitExceededExceptionТаким образом, счетчики предназначены для обработки небольшого количества суммарных сигналов, а не для подсчета количества нажатий каждой клавиши.

Пример счетчиков

Пример класса MapClass с счетчиками для подсчета количества пропущенных и недопустимых значений. Входной файл данных, используемый в этом руководстве: наш входной набор данных — это CSV-файл SalesJan2009.csv.

public static class MapClass
            extends MapReduceBase
            implements Mapper<LongWritable, Text, Text, Text>
{
    static enum SalesCounters { MISSING, INVALID };
    public void map ( LongWritable key, Text value,
                 OutputCollector<Text, Text> output,
                 Reporter reporter) throws IOException
    {
        
        //Input string is split using ',' and stored in 'fields' array
        String fields[] = value.toString().split(",", -20);
        //Value at 4th index is country. It is stored in 'country' variable
        String country = fields[4];
        
        //Value at 8th index is sales data. It is stored in 'sales' variable
        String sales = fields[8];
      
        if (country.length() == 0) {
            reporter.incrCounter(SalesCounters.MISSING, 1);
        } else if (sales.startsWith("\"")) {
            reporter.incrCounter(SalesCounters.INVALID, 1);
        } else {
            output.collect(new Text(country), new Text(sales + ",1"));
        }
    }
}

Приведённый выше фрагмент кода демонстрирует пример реализации счётчиков в Hadoop MapReduce.

Здесь, Счетчики продаж — это счётчик, определённый с помощью 'перечислениеОн используется для подсчета ОТСУТСТВУЮЩИХ и НЕПРАВИЛЬНЫХ входных записей.

В приведенном фрагменте кода, если 'странаЕсли поле имеет нулевую длину, то его значение отсутствует, и, следовательно, соответствующий счетчик SalesCounters.MISSING увеличивается.

Далее, если 'главнаяЕсли поле начинается с символа ", то запись считается НЕДЕЙСТВИТЕЛЬНОЙ. Это указывается увеличением счетчика SalesCounters.INVALID.

💡 Примечание к API: Приведённый выше фрагмент использует оригинал. org.apache.hadoop.mapred API, где MapReduceBase, Mapper интерфейс, OutputCollector и Reporter отображаются отдельно. Текущий код написан для org.apache.hadoop.mapreduceгде один Context заменяет сборщик и репортер, а счетчик увеличивается на единицу. context.getCounter(SalesCounters.MISSING).increment(1)Концепция прилавка идентична в обоих случаях.

Часто задаваемые вопросы (FAQ)

Выбирайте вариант с маппером, если одна сторона достаточно мала, чтобы храниться в памяти на каждом узле, поскольку он полностью пропускает перемешивание. Выбирайте вариант с редуктором, если обе стороны велики или не отсортированы, и примите дополнительные сетевые затраты.

Модели обучаются на основе истории выполнения заданий, чтобы прогнозировать время выполнения, рекомендовать размеры разделения и количество редукторов, а также выявлять отклонения от значений счетчиков. Они также отмечают задания, у которых счетчики пропущенных записей или неудачных задач выходят за пределы нормального диапазона для данного конвейера.

Copilot создает правдоподобные шаблоны мапперов и редьюсеров, но он свободно смешивает старый пакет mapred с более новым пакетом mapreduce в одном классе, что не компилируется. Исправьте импорты и сигнатуры методов, прежде чем доверять этой логике.

Это механизм, который отправляет меньший по размеру файл на каждый узел перед началом выполнения задач. Затем каждый маппер загружает эту копию в хэш-карту и ищет совпадения локально, что и делает возможным объединение на стороне маппера.

Они выводятся в сводке консоли по завершении задания, отображаются на страницах истории заданий и веб-странице менеджера ресурсов, а также могут быть программно считаны из объекта задания, поэтому драйвер может проверить их и признать выполнение некорректным.

Один объект Context. Он выполняет работу, которую ранее разделяли между собой OutputCollector и Reporter, поэтому вывод и увеличение счетчиков осуществляются через один и тот же дескриптор, переданный в метод map.

В большинстве случаев, связанных с журналистской работой, нет. HiveQL join В результате компиляции в несколько строк кода получается тот же шаблон перемешивания и слияния. Ручная работа оправдана, когда логика слияния не соответствует SQL-запросу.

Hadoop отказывается записывать данные в уже существующий выходной каталог, что защищает завершенные результаты от перезаписи. Для этого сначала рекурсивно удалите каталог или укажите другой путь для вывода при следующем запуске.

Подведем итог этой публикации следующим образом: