Hadoop MapReduce Join & Counter із прикладом

⚡ Розумний підсумок

Об'єднання MapReduce об'єднують два великі набори даних на спільному ключі, або всередині маппера, або всередині редуктора, тоді як лічильники MapReduce збирають статистику про завдання, щоб можна було виміряти погані записи, а не вгадувати їх.

  • 🔘 Основи приєднання: Менший з двох наборів даних розподіляється між кожним вузлом даних і використовується як сторона пошуку.
  • ☑️ З'єднання з боку карти: Вимагає, щоб кожен вхідний параметр був розділений на секції, рівномірно розділений та відсортований за ключем об'єднання перед запуском функції map.
  • З'єднання зі скороченням: Не потребує розбиття, оскільки кожен кортеж, що має спільний ключ з'єднання, потрапляє в один і той самий редуктор.
  • 🧪 Приклад роботи: Файли DeptName.txt та DeptStrength.txt копіюються в HDFS та об'єднуються на Dept_ID за допомогою упакованого JAR-архіву.
  • 🛠️ Типи лічильників: П'ять вбудованих груп лічильників постачаються з кожним завданням, а визначені користувачем лічильники оголошуються як Java перерахування.
  • ⚠️ Використання лічильника: Збільшення лічильника для кожного відсутнього або недійсного запису перетворює проблеми з якістю даних на число у звіті про завдання.

Підручник зі з'єднання та лічильника Hadoop MapReduce з робочим прикладом

Що таке Join в MapReduce?

Операція MapReduce Join використовується для об'єднання двох великих наборів даних. Однак цей процес передбачає написання великої кількості коду для виконання фактичної операції об'єднання. Об'єднання двох наборів даних починається з порівняння розміру кожного набору даних. Якщо один набір даних менший порівняно з іншим набором даних, то менший набір даних розподіляється між кожним вузлом даних у кластері.

Після приєднання MapReduce розподілений, або Mapper, або Reducer використовує менший набір даних для пошуку відповідних записів з великого набору даних, а потім об'єднує ці записи для формування вихідних записів.

Типи приєднання

Залежно від місця, де виконується фактичне з'єднання, з'єднання в Hadoop класифікуються на два типи.

  1. З'єднання з боку карти — Коли об'єднання виконується маппером, це називається об'єднанням на стороні мапи. У цьому типі об'єднання виконується до того, як дані фактично споживаються функцією мапи. Обов'язково, щоб вхідні дані для кожної мапи були у формі розділу та були відсортовані. Крім того, має бути однакова кількість розділів, і вони повинні бути відсортовані за ключем об'єднання.
  2. З'єднання зі скороченням — Коли об'єднання виконується редуктором, це називається об'єднанням на стороні редукції. У цьому об'єднанні немає необхідності мати набір даних у структурованій формі (або розділеному). Тут обробка на стороні карти генерує ключ об'єднання та відповідні кортежі обох таблиць. В результаті цієї обробки всі кортежі з однаковим ключем об'єднання потрапляють до одного редуктора, який потім об'єднує записи з однаковим ключем об'єднання.

Загальний процес об’єднання в Hadoop зображено на схемі нижче.

Блок-схема процесу, що порівнює об'єднання на стороні мапи з об'єднанням на стороні редукції в Hadoop
Типи об’єднань у Hadoop MapReduce

Після того, як два варіанти були зрозумілі, у наступному розділі розглянемо об'єднання на стороні редукції для двох невеликих файлів відділів.

Як об’єднати два набори даних: приклад MapReduce

У двох різних файлах (показано нижче) є два набори даних. Ключ Dept_ID є спільним в обох файлах. Мета полягає в тому, щоб використати MapReduce Join для об'єднання цих файлів.

Перший вхідний файл, що містить ідентифікатори відділів разом із назвами відділів

файл 1
Другий вхідний файл, що містить ідентифікатори відділів разом зі значеннями чисельності відділів

файл 2

Вхідний сигнал: Вхідний набір даних – це текстовий файл DeptName.txt та DeptStrength.txt

Завантажте вхідні файли звідси

Переконайтеся, що у вас є Hadoop встановлено. Перш ніж розпочати роботу з прикладом MapReduce Join, змініть користувача на «hduser» (ідентифікатор, який використовується під час налаштування Hadoop, ви можете переключитися на ідентифікатор користувача, який використовується під час налаштування Hadoop).

su - hduser_

Запит змінюється на обліковий запис Hadoop, як показано нижче.

Термінал після перемикання на обліковий запис hduser за допомогою команди su

Крок 1) Скопіюйте файл zip у вибране вами місце

Завантажений архів MapReduceJoin розміщено у вибраному робочому каталозі

Крок 2) Розпакуйте файл Zip

sudo tar -xvf MapReduceJoin.tar.gz

КолишнійtracНазви файлів ted прокручуються, поки tar розпаковує архів.

Список файлів у консолі, наприкладtracотримано з 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

Альтернативою є використання іншої назви вихідного каталогу.

З’єднання показують, як виглядають дані. Лічильники, які розглядаються далі, показують, як поводилося завдання, яке їх створило.

Що таке лічильник у MapReduce?

Лічильник у MapReduce – це механізм, який використовується для збору та вимірювання статистичної інформації про завдання та події MapReduce. Лічильники зберігають track різних статистичних даних про завдання в MapReduce, таких як кількість виконаних операцій та прогрес операції. Лічильники використовуються для діагностики проблем у MapReduce.

Лічильники Hadoop подібні до розміщення повідомлення журналу в коді карти чи зменшення. Ця інформація може бути корисною для діагностики проблеми в обробці завдань MapReduce.

Зазвичай ці лічильники в Hadoop визначені в програмі (map або reduce) та збільшуються під час виконання, коли відбувається певна подія або умова (специфічна для цього лічильника). Дуже гарним застосуванням лічильників Hadoop є track дійсних та недійсних записів з вхідного набору даних.

Типи лічильників MapReduce

Існує два основних типи лічильників MapReduce

  1. Вбудовані лічильники Hadoop: Є кілька вбудованих лічильників Hadoop, які існують для кожного завдання. Нижче наведено вбудовані групи лічильників-
    • Лічильники завдань MapReduce — Збирає інформацію про завдання (наприклад, кількість вхідних записів) під час його виконання.
    • Лічильники файлової системи — Збирає інформацію, таку як кількість байтів, прочитаних або записаних завданням.
    • Лічильники FileInputFormat — Збирає інформацію про певну кількість байтів, зчитаних через FileInputFormat.
    • Лічильники FileOutputFormat — Збирає інформацію про певну кількість байтів, записаних через FileOutputFormat.
    • Лічильники завдань — Ці лічильники записують статистику для всієї роботи, таку як кількість запущених завдань для роботи.
  2. Користувацькі лічильники: Окрім вбудованих лічильників, користувач може визначати власні лічильники, використовуючи аналогічні функції, що надаються мовами програмування. Наприклад, у Java, «перелік» використовується для визначення користувацьких лічильників.

💡 Примітка до версії: лічильники завдань підтримувалися відділом завданьTracker під MRv1. На YARN ця роль належить 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)Концепція лічильника ідентична в обох.

Поширені запитання

Оберіть варіант з відображенням, коли одна сторона достатньо мала для зберігання в пам'яті на кожному вузлі, оскільки він повністю пропускає перетасовку. Оберіть варіант з редуктором, коли обидві сторони великі або не відсортовані, і погодьтеся з додатковими витратами мережі.

Моделі навчаються на основі історії завдань, щоб прогнозувати час виконання, рекомендувати розміри розділень та кількість скорочень, а також виявляти перекіс значень лічильників. Вони також позначають завдання, лічильники розлитих записів або невдалих завдань яких виходять за межі нормального діапазону для цього конвеєра.

Copilot створює правдоподібні скелети мапперів та редукторів, але він вільно змішує старий пакет mapred з новим пакетом mapreduce в одному класі, що не компілюється. Виправте імпорт та сигнатури методів, перш ніж довіряти логіці.

Це механізм, який надсилає менший файл до кожного вузла перед початком завдань. Потім кожен маппер завантажує цю копію в хеш-карту та шукає збіги локально, що робить можливим об'єднання на стороні маппера.

Вони друкуються у зведенні консолі після завершення завдання, відображаються на веб-сторінках історії завдань та менеджера ресурсів, а також зчитуються програмно з об'єкта завдання, тому драйвер може застосувати до них авторизацію та завершити невдалий запуск.

Один об'єкт Context. Він виконує роботу, яку OutputCollector та Reporter використовували для розподілу між собою, тому вивід записується, а лічильники збільшуються через той самий дескриптор, що передається в метод map.

Для більшості репортажних робіт, № А Приєднання до HiveQL компілюється до того ж шаблону перетасування та злиття за кілька рядків. Рукописні завдання варті того, коли логіка злиття не відповідає SQL-реченням.

Hadoop відмовляється записувати дані у вихідний каталог, який вже існує, що захищає готові результати від перезапису. Спочатку рекурсивно видаліть каталог або передайте інший вихідний шлях під час наступного запуску.

Підсумуйте цей пост за допомогою: