大数据技术原理


注意:考试编程内容占比较大,且涉及的Hbase、Hive、Linux指令远不止文中提及的内容,文中部分仅是基础。请同学们根据书本内容认真复习(不过MapReduce编程只会出一道大编程题,且由于代码量相对较大,题目大概率不会很难)


MapReduce编程

普通WordCount
import java.io.IOException;
import java.util.StringTokenizer;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.IntWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Job;
import org.apache.hadoop.mapreduce.Mapper;
import org.apache.hadoop.mapreduce.Reducer;
import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
import org.apache.hadoop.util.GenericOptionsParser;

public class WordCount {
    public WordCount() {
    }

    public static void main(String[] args) throws Exception {
        Configuration conf = new Configuration();
        String[] otherArgs = (new GenericOptionsParser(conf, args)).getRemainingArgs();
        if (otherArgs.length < 2) {
            System.err.println("Usage: wordcount<in> [<in>...] <out>");
            System.exit(2);
        }

        Job job = Job.getInstance(conf, "word count");
        job.setJarByClass(WordCount.class);
        job.setMapperClass(WordCount.TokenizerMapper.class);
        job.setCombinerClass(WordCount.IntSumReducer.class);
        job.setReducerClass(WordCount.IntSumReducer.class);
        job.setOutputKeyClass(Text.class);
        job.setOutputValueClass(IntWritable.class);

        for (int i = 0; i < otherArgs.length - 1; ++i) {
            FileInputFormat.addInputPath(job, new Path(otherArgs[i]));
        }
        FileOutputFormat.setOutputPath(job, new Path(otherArgs[otherArgs.length - 1]));
        System.exit(job.waitForCompletion(true) ? 0 : 1);
    }

    public static class TokenizerMapper extends Mapper<Object, Text, Text, IntWritable> {
        private static final IntWritable one = new IntWritable(1);
        private Text word = new Text();

        public TokenizerMapper() {
        }

        public void map(Object key, Text value, Context context) 
                throws IOException, InterruptedException {
            StringTokenizer itr = new StringTokenizer(value.toString());
            while (itr.hasMoreTokens()) {
                this.word.set(itr.nextToken());
                context.write(this.word, one);
            }
        }
    }

    public static class IntSumReducer extends Reducer<Text, IntWritable, Text, IntWritable> {
        private IntWritable result = new IntWritable();

        public IntSumReducer() {
        }

        public void reduce(Text key, Iterable<IntWritable> values, Context context) 
                throws IOException, InterruptedException {
            int sum = 0;
            IntWritable val;
            for (Iterator i$ = values.iterator(); i$.hasNext(); sum += val.get()) {
                val = (IntWritable) i$.next();
            }
            this.result.set(sum);
            context.write(key, this.result);
        }
    }
}
/**
 * TokenizerMapper 类是一个 Hadoop Mapper 类的实现,
 * 用于将输入的文本数据拆分为单词,并为每个单词输出 <word, 1> 键值对。
 * 
 * 输入:<行偏移量, 文本行>
 * 输出:<单词, 1>(用于后续的单词计数)
 */
public static class TokenizerMapper extends Mapper<Object, Text, Text, IntWritable> {
    // 作为每个单词的初始计数
    private static final IntWritable one = new IntWritable(1);
    // 用于存储当前处理的单词
    private Text word = new Text();
    public TokenizerMapper() {
    }
    /**
     * map 方法是 Mapper 的核心方法,处理输入的键值对并生成中间结果
     * 
     * @param key 输入的键,这里通常是文本行的偏移量(未使用)
     * @param value 输入的值,这里是文本行内容
     * @param context Hadoop 提供的上下文对象,用于输出中间结果
     * @throws IOException
     * @throws InterruptedException
     */
    public void map(Object key, Text value, Context context) 
            throws IOException, InterruptedException {
        StringTokenizer itr = new StringTokenizer(value.toString());
        while (itr.hasMoreTokens()) { 
            this.word.set(itr.nextToken());  // 获取下一个单词
            // 输出 <单词, 1> 键值对
            context.write(this.word, one);
        }
    }
}
/**
 * IntSumReducer 类是一个 Hadoop Reducer 的实现,
 * 用于汇总相同单词的出现次数(对 Mapper 输出的 <word, 1> 进行求和)
 * 
 * 输入:<单词, [1, 1, ...]>(来自 Mapper 的输出)
 * 输出:<单词, 总次数>
 */
public static class IntSumReducer extends Reducer<Text, IntWritable, Text, IntWritable> {
    private IntWritable result = new IntWritable();
    public IntSumReducer() {
    }
    /**
     * reduce 方法是 Reducer 的核心方法,对相同键的值进行聚合
     * 
     * @param key 输入的键,这里是单词
     * @param values 可迭代的计数值(多个1的集合)
     * @param context Hadoop 上下文对象,用于输出最终结果
     * @throws IOException
     * @throws InterruptedException
     */
    public void reduce(Text key, Iterable<IntWritable> values, Context context) 
            throws IOException, InterruptedException {
        int sum = 0;  // 初始化计数器
        // 遍历所有值(都是1)并累加
        for (IntWritable val : values) {
            sum += val.get();  // 将IntWritable转为int并累加
        }
        this.result.set(sum);
        // 输出 <单词, 总次数>
        context.write(key, this.result);
    }
}

分组WordCount
import java.io.IOException;
import java.util.StringTokenizer;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.IntWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Job;
import org.apache.hadoop.mapreduce.Mapper;
import org.apache.hadoop.mapreduce.Reducer;
import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
import org.apache.hadoop.util.GenericOptionsParser;

public class LetterGroupWordCount {
    public LetterGroupWordCount() {
    }

    public static void main(String[] args) throws Exception {
        Configuration conf = new Configuration();
        String[] otherArgs = new GenericOptionsParser(conf, args).getRemainingArgs();
        if (otherArgs.length < 2) {
            System.err.println("Usage: lettergroupwordcount <in> [<in>...] <out>");
            System.exit(2);
        }

        Job job = Job.getInstance(conf, "letter group word count");
        job.setJarByClass(LetterGroupWordCount.class);
        job.setMapperClass(FirstLetterMapper.class);
        job.setCombinerClass(LetterCountReducer.class);
        job.setReducerClass(LetterCountReducer.class);
        job.setOutputKeyClass(Text.class);
        job.setOutputValueClass(IntWritable.class);

        for (int i = 0; i < otherArgs.length - 1; ++i) {
            FileInputFormat.addInputPath(job, new Path(otherArgs[i]));
        }
        FileOutputFormat.setOutputPath(job, new Path(otherArgs[otherArgs.length - 1]));
        System.exit(job.waitForCompletion(true) ? 0 : 1);
    }

    /**
     * Mapper类:提取单词首字母作为key
     * 输入:<行偏移量, 文本行>
     * 输出:<首字母, 1>
     */
    public static class FirstLetterMapper extends Mapper<Object, Text, Text, IntWritable> {
        private final static IntWritable one = new IntWritable(1);
        private Text firstLetter = new Text();
        public void map(Object key, Text value, Context context) 
                throws IOException, InterruptedException {
            StringTokenizer itr = new StringTokenizer(value.toString());
            while (itr.hasMoreTokens()) {
                String word = itr.nextToken().trim();
                if (!word.isEmpty()) {
                    // 获取首字母并转为小写(不区分大小写)
                    String letter = word.substring(0, 1).toLowerCase();
                    firstLetter.set(letter);
                    context.write(firstLetter, one);
                }
            }
        }
    }

    /**
     * Reducer类:统计每个首字母的总出现次数
     * 输入:<首字母, [1, 1, ...]>
     * 输出:<首字母, 总次数>
     */
    public static class LetterCountReducer extends Reducer<Text, IntWritable, Text, IntWritable> {
        private IntWritable result = new IntWritable();
        public void reduce(Text key, Iterable<IntWritable> values, Context context)
                throws IOException, InterruptedException {
            int sum = 0;
            for (IntWritable val : values) {
                sum += val.get(); // 累加相同首字母的计数
            }
            result.set(sum);
            context.write(key, result);
        }
    }
}
public static class FirstLetterMapper extends Mapper<Object, Text, Text, IntWritable> {
    private final static IntWritable one = new IntWritable(1);
    private Text firstLetter = new Text();
    public void map(Object key, Text value, Context context) 
        throws IOException, InterruptedException {
        StringTokenizer itr = new StringTokenizer(value.toString());
        while (itr.hasMoreTokens()) {
            String word = itr.nextToken();
            if (!word.isEmpty()) {
                // 获取首字母并转为小写(不区分大小写)
                String letter = word.substring(0, 1).toLowerCase();
                firstLetter.set(letter);
                context.write(firstLetter, one);
            }
        }
    }
}
public static class LetterCountReducer extends Reducer<Text, IntWritable, Text, IntWritable> {
    private IntWritable result = new IntWritable();
    public void reduce(Text key, Iterable<IntWritable> values, Context context)
        throws IOException, InterruptedException {
        int sum = 0;
        for (IntWritable val : values) {
            sum += val.get(); // 累加相同首字母的计数
        }
        result.set(sum);
        context.write(key, result);
    }
}

倒排索引(不考)
import java.io.IOException;
import java.util.StringTokenizer;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Job;
import org.apache.hadoop.mapreduce.Mapper;
import org.apache.hadoop.mapreduce.Reducer;
import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
import org.apache.hadoop.util.GenericOptionsParser;

public class InvertedIndex {
    public InvertedIndex() {
    }

    public static void main(String[] args) throws Exception {
        Configuration conf = new Configuration();
        String[] otherArgs = new GenericOptionsParser(conf, args).getRemainingArgs();
        if (otherArgs.length < 2) {
            System.err.println("Usage: invertedindex <in> [<in>...] <out>");
            System.exit(2);
        }

        Job job = Job.getInstance(conf, "inverted index");
        job.setJarByClass(InvertedIndex.class);
        job.setMapperClass(InvertedIndex.TokenizerMapper.class);
        job.setCombinerClass(InvertedIndex.IntSumReducer.class);
        job.setReducerClass(InvertedIndex.InvertedIndexReducer.class);
        job.setOutputKeyClass(Text.class);
        job.setOutputValueClass(Text.class);

        for (int i = 0; i < otherArgs.length - 1; ++i) {
            FileInputFormat.addInputPath(job, new Path(otherArgs[i]));
        }
        FileOutputFormat.setOutputPath(job, new Path(otherArgs[otherArgs.length - 1]));
        System.exit(job.waitForCompletion(true) ? 0 : 1);
    }

    public static class TokenizerMapper extends Mapper<Object, Text, Text, Text> {
        private Text word = new Text();
        private Text documentId = new Text();

        public void map(Object key, Text value, Context context) 
                throws IOException, InterruptedException {
            // 获取文件名作为文档ID
            String fileName = ((FileSplit)context.getInputSplit()).getPath().getName();
            documentId.set(fileName);
            
            StringTokenizer itr = new StringTokenizer(value.toString());
            while (itr.hasMoreTokens()) {
                word.set(itr.nextToken());
                // 输出格式: <word, documentId>
                context.write(word, documentId);
            }
        }
    }

    public static class IntSumReducer extends Reducer<Text, Text, Text, Text> {
        private Text result = new Text();

        public void reduce(Text key, Iterable<Text> values, Context context) 
                throws IOException, InterruptedException {
            // Combiner阶段: 统计同一文档中单词出现的次数
            HashMap<String, Integer> docCount = new HashMap<>();
            for (Text val : values) {
                String docId = val.toString();
                docCount.put(docId, docCount.getOrDefault(docId, 0) + 1);
            }
            
            StringBuilder sb = new StringBuilder();
            for (Map.Entry<String, Integer> entry : docCount.entrySet()) {
                if (sb.length() > 0) sb.append(";");
                sb.append(entry.getKey()).append(":").append(entry.getValue());
            }
            
            result.set(sb.toString());
            context.write(key, result);
        }
    }

    public static class InvertedIndexReducer extends Reducer<Text, Text, Text, Text> {
        private Text result = new Text();

        public void reduce(Text key, Iterable<Text> values, Context context) 
                throws IOException, InterruptedException {
            // 合并来自不同Mapper的结果
            HashMap<String, Integer> docCount = new HashMap<>();
            for (Text val : values) {
                String[] parts = val.toString().split(";");
                for (String part : parts) {
                    String[] docInfo = part.split(":");
                    String docId = docInfo[0];
                    int count = Integer.parseInt(docInfo[1]);
                    docCount.put(docId, docCount.getOrDefault(docId, 0) + count);
                }
            }
            
            // 构建最终倒排索引格式
            StringBuilder sb = new StringBuilder();
            for (Map.Entry<String, Integer> entry : docCount.entrySet()) {
                if (sb.length() > 0) sb.append(", ");
                sb.append(entry.getKey()).append(":").append(entry.getValue());
            }
            
            result.set(sb.toString());
            context.write(key, result);
        }
    }
}
public static class TonkenizerMapper extends Mapper<Object, Text, Text, Text> {
    // extends Mapper<Object, Text, Text, Text> 的尖括号 < > 中的类型参数定义了 Mapper 类的输入/输出键值对的类型
    private Text word = new Text();
    private Text decumentId = new Text();
    
    public void map(Object key, Text value, Context context) throws IOException, InterruptedException {
        // 获取文件名作为文档ID
        String filename = ((FileSplit)context.getInputSplit()).getPath().getName();
        documentId.set(fileName);
        StringTokenizer itr = new StringTokenizer(value.toString());
        while(itr.hasMoreToken()){
            context.write(word, documentId);
        }
    }
}
public static class IntSumReducer extends Reducer<Text, Text, Text, Text> {
    private Text result = new Text();
    
    public void reduce(Text key, Iterable<Text> values, Context context) throws IOException,, InterruptedException {
        // 统计统一文档中单词出现的次数
        HashMap<String, Integer> dotCount = new HashMap<>();
        for (Text val : values) {  // wc这个遍历写法好啊
            String docId = val.toString();
            docCount.put(docId, docCount.getOrDefault(docId, 0) + 1);
        }
        StringBuilder sb = new StringBuilder;
        for(Map.Entry<String, Integer> entry : doCount.entrySet()) {
            if (sb.length() > 0) sb.append(";");
            sb.append(entry.getKey()).append(":").append(entry.getValue());
        }
        retult.set(sb.toString());
        context.write(key, result);
    }
}
public static class InvertedIndexReducer extends Reducer<Text, Text, Text, Text> {
    private Text result = new Text();

    public void reduce(Text key, Iterable<Text> values, Context context) 
        throws IOException, InterruptedException {
        // 合并来自不同Mapper的结果
        HashMap<String, Integer> docCount = new HashMap<>();
        for (Text val : values) {
            String[] parts = val.toString().split(";");
            for (String part : parts) {
                String[] docInfo = part.split(":");
                String docId = docInfo[0];
                int count = Integer.parseInt(docInfo[1]);
                docCount.put(docId, docCount.getOrDefault(docId, 0) + count);
            }
        }
        // 构建最终倒排索引格式
        StringBuilder sb = new StringBuilder();
        for (Map.Entry<String, Integer> entry : docCount.entrySet()) {
            if (sb.length() > 0) sb.append(", ");
            sb.append(entry.getKey()).append(":").append(entry.getValue());
        }
        result.set(sb.toString());
        context.write(key, result);
    }
}

Linux基本指令
  • cd:切换目录。cd /path/to/directory 切换到指定目录;cd ~ 切换到主目录;cd .. 切换到上一级目录。
  • ls:列出目录内容。ls -l 显示详细信息;ls -a 显示隐藏文件;ls -h 以易读格式显示文件大小。
  • mkdir:创建目录。mkdir directory_name 创建新目录;mkdir -p path/to/directory 创建多级目录。
  • cp:复制文件或目录。cp source destination 复制文件;cp -r source_directory destination_directory 递归复制目录。
  • mv:移动或重命名文件。mv source destination 移动文件;mv old_name new_name 重命名文件。
  • rm:删除文件或目录。rm file 删除文件;rm -r directory 递归删除目录;rm -i file 提示确认删除。
  • cat:显示文件内容。cat file 显示文件内容;cat file1 file2 > output 合并文件内容。
  • tac:反向显示文件内容。tac file 从最后一行开始显示文件内容。
  • more:分页显示文件内容。more file 分页显示文件内容。
  • head:显示文件的前几行。head file 显示前10行;head -n 5 file 显示前5行。
  • tail:显示文件的最后几行。tail file 显示最后10行;tail -n 5 file 显示最后5行;tail -f file 动态显示新增内容。
  • touch:创建空文件或更新时间戳。touch file 创建空文件;touch -t timestamp file 设置时间戳。
  • chown:更改文件所有者。chown user:group file 更改所有者;chown -R user:group directory 递归更改。
  • find:查找文件或目录。find /path -name "pattern" 按名称查找;find . -type f 查找文件;find . -mtime -1 查找最近修改的文件。
  • tar:归档文件。tar -cvf archive.tar file 创建归档;tar -xvf archive.tar 解压归档;tar -czvf archive.tar.gz file 创建压缩归档。
  • grep:搜索文本内容。grep "pattern" file 搜索文本;grep -i "pattern" file 忽略大小写;grep -r "pattern" /path 递归搜索。

Hbase Shell 常用指令

  • 建表 htest,包含列族 infodata

    create 'htest', 'info', 'data'
    
  • 查看数据库有哪些表

    list
    
  • 查看表的定义

    describe 'htest'
    
  • 修改 dataVERSIONS 为 5

    alter 'htest', {NAME => 'data', VERSIONS => 5}
    # NAME => 'data':指定要修改的列族名为 data。
    # VERSIONS => 5:设置列族 data 的 VERSIONS 参数为 5。VERSIONS 参数表示每个单元格可以存储的最大版本数。
    
  • 插入记录

    put 'htest', 'rk01', 'info:name', 'Alice'
    put 'htest', 'rk01', 'info:age', '25'
    put 'htest', 'rk01', 'info:gender', 'Female'
    put 'htest', 'rk01', 'data:java', '85'
    put 'htest', 'rk01', 'data:hadoop', '90'
    put 'htest', 'rk01', 'data:csharp', '78'
    
    put 'htest', 'rk02', 'info:name', 'Bob'
    put 'htest', 'rk02', 'info:age', '30'
    put 'htest', 'rk02', 'info:gender', 'Male'
    put 'htest', 'rk02', 'data:java', '75'
    put 'htest', 'rk02', 'data:hadoop', '88'
    put 'htest', 'rk02', 'data:csharp', '92'
    
    put 'htest', 'rk03', 'info:name', 'Charlie'
    put 'htest', 'rk03', 'info:age', '22'
    put 'htest', 'rk03', 'info:gender', 'Male'
    put 'htest', 'rk03', 'data:java', '88'
    put 'htest', 'rk03', 'data:hadoop', '76'
    put 'htest', 'rk03', 'data:csharp', '85'
    
  • 输入 rk01hadoop 课程成绩 5 次以上

    put 'htest', 'rk01', 'data:hadoop', '90', 1  # 最后的数字是时间戳(版本)
    put 'htest', 'rk01', 'data:hadoop', '92', 2
    put 'htest', 'rk01', 'data:hadoop', '95', 3
    put 'htest', 'rk01', 'data:hadoop', '88', 4
    put 'htest', 'rk01', 'data:hadoop', '91', 5
    
  • get 方法查看 rk01 的成绩

    get 'htest', 'rk01', {COLUMN => 'data:hadoop'}  
    
  • get 方法查看 rk01 的所有版本的成绩

    get 'htest', 'rk01', {COLUMN => 'data:hadoop', VERSIONS => 5}
    
  • get 方法查看 rk01 的 3 个版本的成绩

    get 'htest', 'rk01', {COLUMN => 'data:hadoop', VERSIONS => 3}
    
  • scan 方法查询

    • 查询所有记录:

      scan 'htest'
      
    • 查询单个 rowkey,如 rk01

      scan 'htest', {STARTROW => 'rk01', ENDROW => 'rk01'}
      
  • 删除某个列,如 java

    delete 'htest', 'rk01', 'data:java'
    
  • 删除某个行,如 student 表中 95001 行

    deleteall 'student', '95001'
    
  • 统计行数

    count 'htest'
    
  • 使用表的别名简化输入

    a = get_table 'htest'
    a.get 'rk01'
    a.scan
    
  • 表的删除

    • 禁用表:

      disable 'htest'
      
    • 删除表:

      drop 'htest'
      

Hive 基本指令

  • 创建数据库:CREATE DATABASE mydb;
  • 切换数据库:USE mydb;
  • 查看数据库:SHOW DATABASES;
  • 删除数据库:DROP DATABASE mydb;
  • 创建表:CREATE TABLE mytable (id INT, name STRING, age INT);
  • 查看表结构:DESCRIBE mytable;
  • 查看表:SHOW TABLES;
  • 删除表:DROP TABLE mytable;
  • 插入数据:INSERT INTO mytable VALUES (1, 'Alice', 25);
  • 查询数据:SELECT * FROM mytable;
  • 查询特定列:SELECT name, age FROM mytable;
  • 查询特定条件:SELECT * FROM mytable WHERE age > 20;
  • 更新数据:UPDATE mytable SET age = 26 WHERE name = 'Alice';
  • 删除数据:DELETE FROM mytable WHERE name = 'Alice';
  • 创建分区表:CREATE TABLE mypartitionedtable (id INT, name STRING) PARTITIONED BY (year INT, month INT);
  • 插入分区数据:INSERT INTO mypartitionedtable PARTITION (year=2023, month=10) VALUES (1, 'Alice');
  • 查询分区数据:SELECT * FROM mypartitionedtable WHERE year = 2023 AND month = 10;
  • 创建视图:CREATE VIEW myview AS SELECT name, age FROM mytable WHERE age > 20;
  • 查询视图:SELECT * FROM myview;
  • 删除视图:DROP VIEW myview;
  • 从文件导入数据:LOAD DATA INPATH '/path/to/file' INTO TABLE mytable;
  • 将数据导出到文件:INSERT OVERWRITE DIRECTORY '/path/to/output' SELECT * FROM mytable;
  • 使用函数:SELECT UPPER(name) FROM mytable;
  • 聚合函数:
    • SELECT COUNT(*) FROM mytable;
    • SELECT SUM(age) FROM mytable;
    • SELECT AVG(age) FROM mytable;
  • 分组和排序:
    • SELECT age, COUNT(*) FROM mytable GROUP BY age;
    • SELECT name FROM mytable ORDER BY age DESC;
  • 创建索引:CREATE INDEX myindex ON TABLE mytable (name) AS 'org.apache.hadoop.hive.ql.index.compact.CompactIndexHandler';

Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐