编写控制台游戏程序

现在让我们使用所学的知识完成一个游戏程序。这里我们将不使用任何图形界面,而是制作一个简单的、控制台上运行的字符界面的游戏。

需要注意Windows的控制台程序和Linux的控制台程序需要使用各自不同的方法来实现诸如光标移动,颜色设置等操作,下面分别讲解。

清屏

在Windows下控制台窗口的控制是基于win32 api, 就是那些在cmd下可以执行的命令, 使用之前需要引入头文件windows.h

例如Windows下清除屏幕:

system("cls");

而在Linux下是通过Shell命令,只需要引入stdlib.h即可,比如要实现清除屏幕:

system("clear);

休眠

Window的控制台中休眠,可以调用windows.h中的Sleep(),比如Sleep(1000)函数,这里1000为毫秒数,功能为延时1s后程序向下运行,下面是一个3秒倒计时程序:

for(int i = 3; i >= 1; i--) {
  printf("%d\n", i);
  Sleep(1000);
}

Linux的控制台休眠可以使用stdlib.h中的system()调用Shell命令sleep n, 这里n代表秒数。

Linux下程序需要这样改写:

for(int i = 3; i >= 1; i--) {
  printf("%d\n", i);
  system("sleep 1");
}

windows控制台窗口操作API

Windows.h中定义的用于控制台窗口操作的API函数如下:

  • GetConsoleScreenBufferInfo 获取控制台窗口信息
  • GetConsoleTitle 获取控制台窗口标题
  • ScrollConsoleScreenBuffer 在缓冲区中移动数据块
  • SetConsoleScreenBufferSize 更改指定缓冲区大小
  • SetConsoleTitle 设置控制台窗口标题
  • SetConsoleWindowInfo 设置控制台窗口信息
#include <windows.h>
SetConsoleTitle("hello world!"); // 设置

由于Linux中往往不安装图形界面的,因此一般也不需要控制控制台窗口,这里就不介绍了。

控制台初始化

#include <iostream>
#include <windows.h>
using namespace std;

int main()
{
    //设置控制台窗口标题
    SetConsoleTitle("hello world!");

    //获取控制台窗口信息;
    //GetConsoleScreenBufferInfo(HANDLE hConsoleOutput, CONSOLE_SCREEN_BUFFER_INFO *bInfo)
    //第一个hConsoleOutput参数(标准控制句柄)通过GetStdHandle()函数返回值获得
    //第二个参数CONSOLE_SCREEN_BUFFER_INFO 保存控制台信息结构体指针
        /*数据成员如下:
        {
            COORD dwSize; // 缓冲区大小
            COORD dwCursorPosition; //当前光标位置
            WORD wAttributes; //字符属性
            SMALL_RECT srWindow; //当前窗口显示的大小和位置
            COORD dwMaximumWindowSize; //最大的窗口缓冲区大小
        }
        */
    HANDLE hOutput = GetStdHandle(STD_OUTPUT_HANDLE);
    CONSOLE_SCREEN_BUFFER_INFO bInfo;
    GetConsoleScreenBufferInfo(hOutput, &bInfo);
    cout << "窗口缓冲区大小:" << bInfo.dwSize.X << ", " << bInfo.dwSize.Y << endl;
    cout << "窗口坐标位置:" << bInfo.srWindow.Left << ", " << bInfo.srWindow.Top
         << ", "<< bInfo.srWindow.Right << ", " << bInfo.srWindow.Bottom << endl;

    //设置显示区域坐标
    //SetConsoleWindowInfo(HANDLE, BOOL, SMALL_RECT *);
    SMALL_RECT rc = {0,0, 79, 24}; // 坐标位置结构体初始化
    SetConsoleWindowInfo(hOutput,true ,&rc);
    cout << "窗口显示坐标位置:" << bInfo.srWindow.Left << ", " << bInfo.srWindow.Top
         << ", "<< bInfo.srWindow.Right << ", " << bInfo.srWindow.Bottom << endl;

    //更改指定缓冲区大小
    //SetConsoleScreenBufferSize(HANDLE hConsoleOutput, COORD dwSize)
    //COORD为一个数据结构体
    COORD dSiz = {80, 25};
    SetConsoleScreenBufferSize(hOutput, dSiz);
    cout << "改变后大小:" << dSiz.X << ", " << dSiz.Y << endl;

    //获取控制台窗口标题
    //GetConsoleTitle(LPTSTR lpConsoleTitle, DWORD nSize)
    //lpConsoleTitle为指向一个缓冲区指针以接收包含标题的字符串;nSize由lpConsoleTitle指向的缓冲区大小
    char cTitle[255];
    GetConsoleTitleA(cTitle, 255);
    cout << "窗口标题:" << cTitle << endl;

    // 关闭标准输出设备句柄
    CloseHandle(hOut);
    return 0;
}

设置文本颜色和光标移动控制

#include <iostream>
#include <windows.h>
using namespace std;
int main()
{
    /*设置文本属性
    BOOL SetConsoleTextAttribute(HANDLE hConsoleOutput, WORD wAttributes);//句柄, 文本属性*/

    HANDLE hOut = GetStdHandle(STD_OUTPUT_HANDLE); // 获取标准输出设备句柄
    WORD wr1 = 0xfa;//定义颜色属性;第一位为背景色,第二位为前景色
    SetConsoleTextAttribute(hOut, wr1);
    cout << "hello world!" << endl;

    WORD wr2 = FOREGROUND_RED | FOREGROUND_INTENSITY;//方法二用系统宏定义颜色属性
    SetConsoleTextAttribute(hOut, wr2);
    cout << "hello world!" << endl;

    /*移动文本位置位置
    BOOL ScrollConsoleScreenBuffer(HANDLE hConsoleOutput, CONST SMALL_RECT* lpScrollRectangle, CONST SMALL_RECT* lpClipRectangle,
                                   COORD dwDestinationOrigin,CONST CHAR_INFO* lpFill);
                                  // 句柄// 裁剪区域// 目标区域 // 新的位置// 填充字符*/
    //输出文本
    SetConsoleTextAttribute(hOut, 0x0f);
    cout << "01010101010101010101010101010" << endl;
    cout << "23232323232323232323232323232" << endl;
    cout << "45454545454545454545454545454" << endl;
    cout << "67676767676767676767676767676" << endl;

    SMALL_RECT CutScr = {1, 2, 10, 4}; //裁剪区域与目标区域的集合行成剪切区域
    SMALL_RECT PasScr = {7, 2, 11, 9}; //可以是NULL,即全区域
    COORD pos = {1, 8};     //起点坐标,与裁剪区域长宽构成的区域再与目标区域的集合为粘贴区

    //定义填充字符的各个参数及属性
    SetConsoleTextAttribute(hOut, 0x1);
    CONSOLE_SCREEN_BUFFER_INFO Intsrc;
    GetConsoleScreenBufferInfo(hOut, &Intsrc);
    CHAR_INFO chFill = {'A',  Intsrc.wAttributes}; //定义剪切区域填充字符
    ScrollConsoleScreenBuffer(hOut, &CutScr, &PasScr, pos, &chFill); //移动文本

    CloseHandle(hOut); // 关闭标准输出设备句柄
    return 0;
}

WORD文本属性预定义宏:(可以直接用16进制表示,WORD w = 0xf0;前一位表示背景色,后一位代表前景色)

FOREGROUND_BLUE 蓝色
FOREGROUND_GREEN 绿色
FOREGROUND_RED 红色
FOREGROUND_INTENSITY 加强
BACKGROUND_BLUE 蓝色背景
BACKGROUND_GREEN 绿色背景
BACKGROUND_RED 红色背景
BACKGROUND_INTENSITY 背景色加强
COMMON_LVB_REVERSE_VIDEO 反色

当前文本属性信息可通过调用函数
GetConsoleScreenBufferInfo后,在CONSOLESCREEN BUFFER_INFO结构成员wAttributes中得到。
在指定位置处写属性

BOOL WriteConsoleOutputAttribute(HANDLE hConsoleOutput, CONST WORD *lpAttribute, DWORD nLength, 
                                COORD dwWriteCoord, LPDWORD lpNumberOfAttrsWritten);
                                //句柄, 属性, 个数, 起始位置, 已写个数*/

填充指定数据的字符

BOOL FillConsoleOutputCharacter(HANDLE hConsoleOutput, TCHAR cCharacter,DWORD nLength, 
                               COORD dwWriteCoord, LPDWORD lpNumberOfCharsWritten);
                               // 句柄, 字符, 字符个数, 起始位置, 已写个数*/

在当前光标位置处插入指定数量的字符

BOOL WriteConsole(HANDLE hConsoleOutput, CONST VOID *lpBuffer, DWORD nNumberOfCharsToWrite,
                 LPDWORD lpNumberOfCharsWritten,LPVOID lpReserved);
                 //句柄, 字符串, 字符个数, 已写个数, 保留*/

向指定区域写带属性的字符

BOOL WriteConsoleOutput(HANDLE hConsoleOutput, CONST CHAR_INFO *lpBuffer, COORD dwBufferSize,
                       COORD dwBufferCoord,PSMALL_RECT lpWriteRegion );
                       // 句柄 // 字符数据区// 数据区大小// 起始坐标// 要写的区域*/

在指定位置处插入指定数量的字符

BOOL WriteConsoleOutputCharacter(HANDLE hConsoleOutput, LPCTSTR lpCharacter, DWORD nLength,
                                COORD dwWriteCoord, LPDWORD lpNumberOfCharsWritten);
                                // 句柄// 字符串// 字符个数// 起始位置// 已写个数*/

填充字符属性

BOOL FillConsoleOutputAttribute(HANDLE hConsoleOutput, WORD wAttribute,DWORD nLength,
                               COORD dwWriteCoord, LPDWORD lpNumberOfAttrsWritten);
                               //句柄, 文本属性, 个数, 开始位置, 返回填充的个数*/

设置代码页,代码页是字符集编码的别名,也有人称"内码表"。

SetConsoleOutputCP(437);
//如(简体中文) 设置成936

光标操作控制

#include <iostream>
#include <windows.h>
using namespace std;
int main()
{
    cout << "hello world!" << endl;

    //设置光标位置
    //SetConsoleCursorPosition(HANDLE hConsoleOutput,COORD dwCursorPosition);
    //设置光标信息
    //BOOL SetConsoleCursorInfo(HANDLE hConsoleOutput, PCONST PCONSOLE_CURSOR_INFO lpConsoleCursorInfo);
    //获取光标信息
    //BOOL GetConsoleCursorInfo(HANDLE hConsoleOutput,  PCONSOLE_CURSOR_INFO lpConsoleCursorInfo);
    //参数1:句柄;参数2:CONSOLE_CURSOR_INFO结构体{DWORD dwSize;(光标大小取值1-100)BOOL bVisible;(是否可见)}

    Sleep(2000);//延时函数
    HANDLE hOut = GetStdHandle(STD_OUTPUT_HANDLE);
    COORD w = {0, 0};
    SetConsoleCursorPosition(hOut, w);
    CONSOLE_CURSOR_INFO cursorInfo = {1, FALSE};
    Sleep(2000);//延时函数
    SetConsoleCursorInfo(hOut, &cursorInfo);
    CloseHandle(hOut); // 关闭标准输出设备句柄
    return 0;
}

键盘操作控制

#include <iostream>
#include <windows.h>
#include <conio.h>
using namespace std;
HANDLE hOut;
//清除函数
void cle(COORD ClPos)
{
    SetConsoleCursorPosition(hOut, ClPos);
    cout << "            " << endl;
}
//打印函数
void prin(COORD PrPos)
{
    SetConsoleCursorPosition(hOut, PrPos);
    cout << "hello world!" << endl;
}
//移动函数
void Move(COORD *MoPos, int key)
{
    switch(key)
    {
    case 72: MoPos->Y--;break;
    case 75: MoPos->X--;break;
    case 77: MoPos->X++;break;
    case 80: MoPos->Y++;break;
    default: break;
    }
}

int main()
{
    cout << "用方向键移动下行输出内容" << endl;
    hOut = GetStdHandle(STD_OUTPUT_HANDLE);//取句柄
    COORD CrPos = {0, 1};//保存光标信息
    prin(CrPos);//打印
    //等待键按下
    while(1)
    {
        if(kbhit())
        {
            cle(CrPos);//清除原有输出
            Move(&CrPos, getch());
            prin(CrPos);
        }
    }
    return 0;
}

可以用方向键任意移动hello world!

注意区分:

getch()getche()getcher()函数

Linux下控制台操作替代办法

字体颜色

在 Linux 下若想输出 类似与 Windows 下的多颜色字体如何做呢?本文就来介绍实现的方法。
首先,来看下 在Linux 下颜色的表示

注意自定义的配置需在写在\033[和m之间

\033[22;30m - black
\033[22;31m - red
\033[22;32m - green
\033[22;33m - brown
\033[22;34m - blue
\033[22;35m - magenta
\033[22;36m - cyan
\033[22;37m - gray
\033[01;30m - dark gray
\033[01;31m - light red
\033[01;32m - light green
\033[01;33m - yellow
\033[01;34m - light blue
\033[01;35m - light magenta
\033[01;36m - light cyan
\033[01;37m - white

可以看出都是使用的转义字体来实现的。
比如:
Linux 终端输入:echo -e "\033[35;1m Shocking \033[0m"
C代码: printf("\033[34mThis is blue.\033[0m\n");
是不是出错不同的颜色了,记得最后要 "\033[0m" 关闭所有属性,这样又回到了系统默认的颜色了。

一个定义宏的办法如下:([0;是用来清除之前或之后的设置)

#define COLOR(msg, code) "\033[0;1;" #code "m" msg "\033[0m"
#define RED(msg)    COLOR(msg, 31)
#define GREEN(msg)  COLOR(msg, 32)
#define YELLOW(msg) COLOR(msg, 33)
#define BLUE(msg)   COLOR(msg, 34)

光标控制

Linux 下终端 C 语言控制光标的技巧

// 清除屏幕
#define CLEAR() printf("\033[2J")

// 上移光标
#define MOVEUP(x) printf("\033[%dA", (x))

// 下移光标
#define MOVEDOWN(x) printf("\033[%dB", (x))

// 左移光标
#define MOVELEFT(y) printf("\033[%dD", (y))

// 右移光标
#define MOVERIGHT(y) printf("\033[%dC",(y))

// 定位光标

#define MOVETO(x,y) printf("\033[%d;%dH", (x), (y))

// 光标复位

#define RESET_CURSOR() printf("\033[H")

// 隐藏光标

#define HIDE_CURSOR() printf("\033[?25l")

// 显示光标

#define SHOW_CURSOR() printf("\033[?25h")

//反显

#define HIGHT_LIGHT() printf("\033[7m")

#define UN_HIGHT_LIGHT() printf("\033[27m")

模拟按键监听

windows平台接受字符并不回显可以调用<conio.h>库的getch函数, 但是这是依赖于windows的BIOS。
linux下可以使用命令stty -echo来关闭回显,当接受字符后,使用stty echo命令恢复设置即可.

linux下如何向windows的getch或者getche一样,不需要回车就可以接受字符?
这里可以使用命令raw开启一次接受一个字符的模式,接受字符后,再次使用cooked(回车之后一锅端模式)模式即可

#include <stdio.h>
#include <stdlib.h>
#define ENTER_KEY 13
#define ESCAPE_KEY 27

char getch()
{
    char ch;
    system("stty -echo raw");
    ch = getchar();
    system("stty echo cooked");
    return ch;
}

int main(int argc, char *argv[])
{
    char ch;

    system("clear\n");
    while (1)
    {
        printf("Press Enter to continue, ESC to exit\n");
        ch = getch();
        if (ch == ESCAPE_KEY)
            break;
        if (ch == ENTER_KEY)
        {
            printf("Input a character\n");
            ch = getch();
            printf("=%c\n", ch);
        }
    }

    return 0;
}

Hangman 猜词游戏

我们将从构建经典的猜词游戏,需要有两个玩家,一个玩家负责出题,另一个玩家负责猜。如果你从来没有听过这个游戏,

规则

  1. 在给定的选词范围内,比如"动物"或""国家",玩家一选择一个秘密单词并根据字母数画相应数量的短横,即“”,短横用来表示单词中的每个字母。比如秘密单词是"cat"则可以用" "表示。

  2. 游戏开始后,玩家二猜一个字母,玩家一判断这个字符是否在秘密单词中出现,如果出现则将对应位置的短横替换成正确的字符。比如猜的是a,则表示成_ _ a

  3. 如果玩家二猜错了字母,玩家一会逐步从上到下的完成一副小人上吊的图像,如果玩家2在完整的图像被画出来之前猜出来秘密单词的所有字母他就赢了,否则当完整的图像画出来之后,小人的脚离开地面就预示着玩家二输了。

    第1次猜错

    ________   

    第2次猜错

    ________   
    |      |   

    第3次猜错

    ________   
    |      |   
    |      0   

    第4次猜错

    ________   
    |      |   
    |      0   
    |     /|\  

    第5次猜错

    ________   
    |      |   
    |      0   
    |     /|\  
    |     / \  

    第6次猜错,小人脚离地,玩家二输了,游戏结束

    ________   
    |      |   
    |      0   
    |     /|\  
    |     / \  
    |          
  4. 如果玩家在小人脚离地之前猜出了所有字母,例如 c a t.则玩家二获胜。

注意:以下代码是在linux环境中运行:

游戏的初始化
// 初始化词库
char words[5][100] = {"china", "japan", "india", "korea", "syria"};
// 初始化词库的数量
int words_num = sizeof(words) / sizeof(words[0]);
// 初始化要猜的词
char secret_word[100] = "";
// 初始化要猜的词的字母数
int secret_word_len = 0;
// 显示的题板, 未猜出的字母用_代替
// 打印的上吊小人的画像
char stages[][15] = {"  ________    ",   // stage 1
                     "  |      |    ",   // stage 2
                     "  |      0    ",   // stage 3
                     "  |     /|\\  ",   // stage 4
                     "  |     / \\  ",   // stage 5
                     "  |________   "};  // stage 6
int stage_num     = sizeof(stages) / sizeof(stages[0]);
// 初始化猜错的次数
int wrong_guess = 0;
// 初始化是否赢得游戏
bool win = false;

// FUNCTION PROTOTYPES
void PrintBoard(int stage);
void StartGame();
char getch();

int main(int argc, char *argv[]) {

    // init
    srand((unsigned int)time(NULL));
    int random_index = rand() % words_num;
    strcpy(secret_word, words[random_index]);
    secret_word_len = strlen(secret_word);

    StartGame();
    return 0;
}
玩法部分

游戏的玩法部分封装到了StartGame函数中,玩法部分通常是个循环


// GAME LOOP
while (wrong_guess < stage_num) {
...
}

判断是否猜中

char *find_pos = strstr(remaining_letters, guess);
if (find_pos != NULL) {
    int pos                = find_pos - remaining_letters;
    remaining_letters[pos] = '$';
    letter_board[pos]      = guess[0];
}
else {
    wrong_guess++;
    PrintBoard(wrong_guess);
    system("sleep 1");
}

如果输了就打印上吊图像


void PrintBoard(int stage) {
    for (int i = 0; i < stage_num; i++) {
        if (i < stage)
            printf("%s\n", stages[i]);
        else
            puts("");
    }
}

练习

你也可以试着完成一个属于你自己的控制台小游戏, 多花些时间做些类似的练习可以帮助你提高编程能力,同时也能让你对编程保持兴趣。

我还记得在DOS时期(95年左右)我第一个使用QuickBasic语言制作的猜飞机头的游戏,虽然是简单的游戏但是最终被我制作的越来越复杂,比如增加关卡,添加声效,游戏记录的保存等,可惜是存储在软盘中现在已经遗失了。

  • 三子棋或者五子棋游戏

  • 根据心情选歌

  • 根据品牌打印广告语。

寻求帮助

如果你写代码时卡壳,大多数的问题通过百度都可以解决,不行的话就CSDN或者知乎。

如果英文水平不赖就科学上网去谷歌搜索,这里要强调的是浏览Stack Exchange这个神奇的网站会很有帮助。有两个版块非常有用。

一个是Stack Overflow,程序员应该没有不知道它的,上面几乎涵盖所有编程可能遇到的问题。即便你的问题是独一无二其他人都没有经历过的,这是一个QA类型的网站,你可以在上面发布问题,会有大佬帮你解答的。

https://www.codementor.io/ama/0926528143/stackoverflow-python-moderator-martijn-pieters-zopatista

http://programmers.stackexchange.com/questions/44177/what-is-the-single-most-effective-thing-you-did-to-improve-your-programming-skil

另外一个有用的板块叫Code Review, 你只要发布你的代码就会有人给你"指手画脚", 告诉你哪些地方做的好,哪些地方做的不好,如何改进之类的建议。

Views: 36

Maxwell 数据库数据实时采集

1、Maxwell 简介

Maxwell 是一个能实时读取 MySQL 二进制日志文件binlog,并生成 Json格式的消息,作为生产者发送给 Kafka,Kinesis、RabbitMQ、Redis、Google Cloud Pub/Sub、文件或其它平台的应用程序。它的常见应用场景有ETL、维护缓存、收集表级别的dml指标、增量到搜索引擎、数据分区迁移、切库binlog回滚方案等。

Maxwell主要提供了下列功能

    1. 支持SELECT * FROM table的方式进行全量数据初始化。
    1. 支持在主库发生failover后,自动恢复binlog位置,实现断点续传。
    1. 可以对数据进行分区,解决数据倾斜问题,发送到Kafka的数据支持库、表、列等级别的数据分区。
    1. 工作方式是伪装为slave接收binlog events,然后根据schema信息拼装,可以接受ddl、xid、row等event。

2、Mysql Binlog介绍

2.1 Binlog 简介

MySQL中一般有以下几种日志

日志类型 写入日志的信息
错误日志 记录在启动,运行或停止mysqld时遇到的问题
通用查询日志 记录建立的客户端连接和执行的语句
二进制日志 binlog 记录更改数据的语句
中继日志 从服务器 复制 主服务器接收的数据更改
慢查询日志 记录所有执行时间超过 long_query_time 秒的所有查询或不使用索引的查询
DDL日志(元数据日志) 元数据操作由DDL语句执行

在默认情况下,系统仅仅打开错误日志,关闭了其他所有日志,以达到尽可能减少IO损耗提高系统性能的目的,但是在一般稍微重要一点的实际应用场景中,都至少需要打开二进制日志,因为这是MySQL很多存储引擎进行增量备份的基础,也是MySQL实现复制的基本条件

接下来主要介绍二进制日志 binlog。

MySQL 的二进制日志 binlog 可以说是 MySQL 最重要的日志,它记录了所有的 DDLDML 语句(除了数据查询语句select、show等),以事件形式记录,还包含语句所执行的消耗的时间,MySQL的二进制日志是事务安全型的。binlog 的主要目的是复制和恢复

Binlog日志的两个最重要的使用场景

  • MySQL主从复制
    • MySQL Replication在Master端开启binlog,Master把它的二进制日志传递给slaves来达到master-slave数据一致的目的。
  • 数据恢复
    • 通过使用 mysqlbinlog工具来使恢复数据。

2.2 Binlog 的日志格式

记录在二进制日志中的事件的格式取决于二进制记录格式。支持三种格式类型:

  • Statement:基于SQL语句的复制(statement-based replication, SBR)
  • Row:基于行的复制(row-based replication, RBR)
  • Mixed:混合模式复制(mixed-based replication, MBR)

Statement

  • 每一条会修改数据的sql都会记录在binlog中。
  • 优点
    • 不需要记录每一行的变化,减少了binlog日志量,节约了IO, 提高了性能。
  • 缺点
    • 在进行数据同步的过程中有可能出现数据不一致。
    • 比如 update tt set create_date=now(),如果用binlog日志进行恢复,由于执行时间不同可能产生的数据就不同。

Row

  • 它不记录sql语句上下文相关信息,仅保存哪条记录被修改。
  • 优点
    • 保持数据的绝对一致性。因为不管sql是什么,引用了什么函数,它只记录执行后的效果。
  • 缺点
    • 每行数据的修改都会记录,最明显的就是update语句,导致更新多少条数据就会产生多少事件,占用较大空间。

Mixed

  • 从5.1.8版本开始,MySQL提供了Mixed格式,实际上就是Statement与Row的结合。
  • 在Mixed模式下,一般的复制使用Statement模式保存binlog,对于Statement模式无法复制的操作使用Row模式保存binlog, MySQL会根据执行的SQL语句选择日志保存方式(因为statement只有sql,没有数据,无法获取原始的变更日志,所以一般建议为Row模式)。
  • 优点
    • 节省空间,同时兼顾了一定的一致性。
  • 缺点
    • 还有些极个别情况依旧会造成不一致,另外statement和mixed对于需要对binlog的监控的情况都不方便。

3、Mysql 实时数据同步方案对比

  • mysql 数据实时同步可以通过解析mysql的 binlog 的方式来实现,解析binlog可以有多种方式,可以通过canal,或者maxwell等各种方式实现。以下是各种抽取方式的对比介绍。

    mysql实时同步方案对比

  • 其中canal 由 Java开发,分为服务端和客户端,拥有众多的衍生应用,性能稳定,功能强大;canal 需要自己编写客户端来消费canal解析到的数据。

  • Maxwell相对于canal的优势是使用简单,Maxwell比Canal更加轻量级,它直接将数据变更输出为json字符串,不需要再编写客户端。对于缺乏基础建设,短时间内需要快速迭代的项目和公司比较合适。

  • 另外Maxwell 有一个亮点功能,就是Canal只能抓取最新数据,对已存在的历史数据没有办法处理。而Maxwell有一个bootstrap功能,可以直接引导出完整的历史数据用于初始化,非常好用。

4、开启Mysql的Binlog

  • 1、服务器当中安装mysql(省略)

    • 注意:mysql的版本尽量不要太低,也不要太高,最好使用5.6及以上版本。
  • 2、添加mysql普通用户maxwell

    • 为mysql添加一个普通用户maxwell,因为maxwell这个软件默认用户使用的是maxwell这个用户。

    • 进入mysql客户端,然后执行以下命令,进行授权

    mysql -uroot -p123456
    • 执行sql语句
    --校验级别最低,只校验密码长度
    mysql> set global validate_password_policy=LOW;
    mysql> set global validate_password_length=6;
    
    --创建maxwell库(启动时候会自动创建,不需手动创建)和用户
    mysql> CREATE USER 'maxwell'@'%' IDENTIFIED BY '123456';
    mysql> GRANT ALL ON maxwell.* TO 'maxwell'@'%';
    mysql> GRANT SELECT, REPLICATION CLIENT, REPLICATION SLAVE on *.* to 'maxwell'@'%'; 
    --刷新权限
    mysql> flush privileges;

    maxwell会自动在MySQL中创建名为maxwell的数据库作为元数据保存使用。

  • 3、修改配置文件 /etc/my.cnf

    • 执行命令 sudo vim /etc/my.cnf, 添加或修改以下三行配置
    #binlog日志名称前缀
    log-bin= /var/lib/mysql/mysql-bin
    
    #binlog日志格式
    binlog-format=ROW
    
    #唯一标识,这个值的区间是:1到(2^32)-1
    server_id=1
  • 4、重启mysql服务

    • 执行如下命令
    sudo service mysqld restart
  • 5、验证binlog是否配置成功

    • 进入mysql客户端,并执行以下命令进行验证
    mysql -uroot -p123456
    mysql> show variables like '%log_bin%';

    image-20210517161330004

  • 6、查看binlog日志文件生成

    • 进入 /var/lib/mysql 目录,查看binlog日志文件.

    image-20210517162315281

5、Maxwell安装部署

  • 1、下载对应版本的安装包

  • 2、上传服务器

  • 3、解压安装包到指定目录

    tar -zxvf maxwell-1.21.1.tar.gz -C /kkb/install/
  • 4、修改maxwell配置文件

    • 进入到安装目录 /kkb/install/maxwell-1.21.1 进行如下操作
    cd /kkb/install/maxwell-1.21.1 
    cp config.properties.example config.properties
    vim config.properties
    • 配置文件config.properties 内容如下:
    # choose where to produce data to
    producer=kafka
    # list of kafka brokers
    kafka.bootstrap.servers=node01:9092,node02:9092,node03:9092
    # mysql login info
    host=node03
    port=3306
    user=maxwell
    password=123456
    # kafka topic to write to
    kafka_topic=maxwell
    • 注意:一定要保证使用maxwell 用户和 123456 密码能够连接上mysql数据库。

6、kafka介绍和使用

6.1 Kafka简介

​ Kafka是最初由Linkedin公司开发,它是一个分布式、可分区、多副本,基于zookeeper协调的分布式日志系统;常见可以用于web/nginx日志、访问日志,消息服务等等。Linkedin于2010年贡献给了Apache基金会并成为顶级开源项目。主要应用场景是:日志收集系统和消息系统

​ Kafka是一个分布式消息队列。具有高性能、持久化、多副本备份、横向扩展能力。生产者往队列里写消息,消费者从队列里取消息进行业务逻辑。Kafka就是一种发布-订阅模式。将消息保存在磁盘中,以顺序读写方式访问磁盘,避免随机读写导致性能瓶颈。

  • 消息(Message)
    • 是指在应用之间传送的数据,消息可以非常简单,比如只包含文本字符串,也可以更复杂,可能包含嵌入对象。
  • 消息队列(Message Queue)
    • 一种应用间的通信方式,消息发送后可以立即返回,通过消息系统来确保信息的可靠传递,消息发布者只管把消息发布到MQ中而不管谁来取,消息使用者只管从MQ中取消息而不管谁发布的,这样发布者和使用者都不用知道对方的存在。

6.2 Kafka特性

  • 高吞吐、低延迟

    kafka 最大的特点就是收发消息非常快,kafka 每秒可以处理几十万条消息,它的最低延迟只有几毫秒。
  • 高伸缩性

    每个主题(topic) 包含多个分区(partition),主题中的分区可以分布在不同的主机(broker)中。
  • 持久性、可靠性

    Kafka 能够允许数据的持久化存储,消息被持久化到磁盘,并支持数据备份防止数据丢失。
  • 容错性

    允许集群中的节点失败,某个节点宕机,Kafka 集群能够正常工作。
  • 高并发

    支持数千个客户端同时读写。

6.3 Kafka集群架构

kafka集群架构

  • producer

    消息生产者,发布消息到Kafka集群的终端或服务。
  • broker

    Kafka集群中包含的服务器,一个borker就表示kafka集群中的一个节点。
  • topic

    每条发布到Kafka集群的消息属于的类别,即Kafka是面向 topic 的。
    更通俗的说Topic就像一个消息队列,生产者可以向其写入消息,消费者可以从中读取消息,一个Topic支持多个生产者或消费者同时订阅它,所以其扩展性很好。
  • partition

    每个 topic 包含一个或多个partition。Kafka分配的单位是partition。
  • replica

    partition的副本,保障 partition 的高可用。
  • consumer

    从Kafka集群中消费消息的终端或服务。
  • consumer group

    每个 consumer 都属于一个 consumer group,每条消息只能被 consumer group 中的一个 Consumer 消费,但可以被多个 consumer group 消费。
  • leader

    每个partition有多个副本,其中有且仅有一个作为Leader,Leader是当前负责数据的读写的partition。 producer 和 consumer 只跟 leader 交互。
  • follower

    Follower跟随Leader,所有写请求都通过Leader路由,数据变更会广播给所有Follower,Follower与Leader保持数据同步。如果Leader失效,则从Follower中选举出一个新的Leader。
  • controller

    知道大家有没有思考过一个问题,就是Kafka集群中某个broker宕机之后,是谁负责感知到他的宕机,以及负责进行Leader Partition的选举?如果你在Kafka集群里新加入了一些机器,此时谁来负责把集群里的数据进行负载均衡的迁移?包括你的Kafka集群的各种元数据,比如说每台机器上有哪些partition,谁是leader,谁是follower,是谁来管理的?如果你要删除一个topic,那么背后的各种partition如何删除,是谁来控制?还有就是比如Kafka集群扩容加入一个新的broker,是谁负责监听这个broker的加入?如果某个broker崩溃了,是谁负责监听这个broker崩溃?这里就需要一个Kafka集群的总控组件,Controller。他负责管理整个Kafka集群范围内的各种东西。
    
  • zookeeper

    (1)   Kafka 通过 zookeeper 来存储集群的meta元数据信息。
    (2)一旦controller所在broker宕机了,此时临时节点消失,集群里其他broker会一直监听这个临时节点,发现临时节点消失了,就争抢再次创建临时节点,保证有一台新的broker会成为controller角色。
  • offset

    • 偏移量
    消费者在对应分区上已经消费的消息数(位置),offset保存的地方跟kafka版本有一定的关系。
    kafka0.8 版本之前offset保存在zookeeper上。
    kafka0.8 版本之后offset保存在kafka集群上。
    它是把消费者消费topic的位置通过kafka集群内部有一个默认的topic,
    名称叫 __consumer_offsets,它默认有50个分区。

6.4 Kafka集群安装部署

  • 1、下载安装包(http://kafka.apache.org

    kafka_2.11-1.1.0.tgz
  • 2、规划安装目录

    /kkb/install
  • 3、上传安装包到服务器中

    通过FTP工具上传安装包到node01服务器上
  • 4、解压安装包到指定规划目录

    tar -zxvf kafka_2.11-1.1.0.tgz -C /kkb/install
  • 5、重命名解压目录

    mv kafka_2.11-1.1.0 kafka
  • 6、修改配置文件

    • 在node01上修改

    • 进入到kafka安装目录下有一个config目录

      • vi server.properties
      #指定kafka对应的broker id ,唯一
      broker.id=0
      #指定数据存放的目录
      log.dirs=/kkb/install/kafka/kafka-logs
      #指定zk地址
      zookeeper.connect=node01:2181,node02:2181,node03:2181
      #指定是否可以删除topic ,默认是false 表示不可以删除
      delete.topic.enable=true
      #指定broker主机名
      host.name=node01
    • 配置kafka环境变量

      • sudo vi /etc/profile
      export KAFKA_HOME=/kkb/install/kafka
      export PATH=$PATH:$KAFKA_HOME/bin
  • 6、分发kafka安装目录到其他节点

    scp -r kafka node02:/kkb/install
    scp -r kafka node03:/kkb/install
    scp /etc/profile node02:/etc
    scp /etc/profile node03:/etc
  • 7、修改node02和node03上的配置

    • node02

    • vi server.properties

      #指定kafka对应的broker id ,唯一
      broker.id=1
      #指定数据存放的目录
      log.dirs=/kkb/install/kafka/kafka-logs
      #指定zk地址
      zookeeper.connect=node01:2181,node02:2181,node03:2181
      #指定是否可以删除topic ,默认是false 表示不可以删除
      delete.topic.enable=true
      #指定broker主机名
      host.name=node02
    • node03

    • vi server.properties

      #指定kafka对应的broker id ,唯一
      broker.id=2
      #指定数据存放的目录
      log.dirs=/kkb/install/kafka/kafka-logs
      #指定zk地址
      zookeeper.connect=node01:2181,node02:2181,node03:2181
      #指定是否可以删除topic ,默认是false 表示不可以删除
      delete.topic.enable=true
      #指定broker主机名
      host.name=node03
  • 8、让每台节点的kafka环境变量生效

    • 在每台服务器执行命令
    source /etc/profile

6.5 kafka集群启动和停止

  • 1、启动kafka集群

    • 先启动zookeeper集群,然后在所有节点如下执行脚本
    nohup kafka-server-start.sh /kkb/install/kafka/config/server.properties >/dev/null 2>&1 &
  • 2、停止kafka集群

    • 所有节点执行关闭kafka脚本
    kafka-server-stop.sh

6.6 kafka命令行的管理使用

  • 1、创建topic

    • 使用 kafka-topics.sh脚本
    kafka-topics.sh --create --partitions 3 --replication-factor 2 --topic test --zookeeper node01:2181,node02:2181,node03:2181
  • 2、查询所有的topic

    • 使用 kafka-topics.sh脚本
    kafka-topics.sh --list --zookeeper node01:2181,node02:2181,node03:2181 
  • 3、查看topic的描述信息

    • 使用 kafka-topics.sh脚本
    kafka-topics.sh --describe --topic test --zookeeper node01:2181,node02:2181,node03:2181  
  • 4、删除topic

    • 使用 kafka-topics.sh脚本
    kafka-topics.sh --delete --topic test --zookeeper node01:2181,node02:2181,node03:2181 
  • 5、模拟生产者写入数据到topic中

    • 使用 kafka-console-producer.sh 脚本
    kafka-console-producer.sh --broker-list node01:9092,node02:9092,node03:9092 --topic test 
  • 6、模拟消费者拉取topic中的数据

    • 使用 kafka-console-consumer.sh 脚本
    kafka-console-consumer.sh --zookeeper node01:2181,node02:2181,node03:2181 --topic test --from-beginning

    或者(推荐)

    kafka-console-consumer.sh --bootstrap-server node01:9092,node02:9092,node03:9092 --topic test --from-beginning

7、Maxwell实时采集mysql表数据到kafka

  • 1、启动kafka集群和zookeeper集群

    • 启动zookeeper集群
    #每台节点执行脚本
    nohup zkServer.sh start >/dev/null  2>&1 &
    • 启动kafka集群
    nohup /kkb/install/kafka/bin/kafka-server-start.sh /kkb/install/kafka/co
    nfig/server.properties > /dev/null 2>&1 &
  • 2、创建topic

    kafka-topics.sh --create --topic maxwell --partitions 3 --replication-factor 2 --zookeeper node01:2181,node02:2181,node03:2181

    (如果虚拟机磁盘容量有效,可以将分区数和复制因子都设置为1.)

  • 3、启动maxwell服务

    /kkb/install/maxwell-1.21.1/bin/maxwell
  • 4、插入数据并进行测试

    • 向mysql表中插入一条数据,并开启kafka的消费者,查看kafka是否能够接收到数据。

    • 向mysql当中创建数据库和数据库表并插入数据

      CREATE DATABASE /*!32312 IF NOT EXISTS*/<code>test_db /*!40100 DEFAULT CHARACTER SET utf8 */;
      
      USE test_db;
      
      /*Table structure for table user */
      
      DROP TABLE IF EXISTS user;
      
      CREATE TABLE user (
      id varchar(10) NOT NULL,
      name varchar(10) DEFAULT NULL,
      age int(11) DEFAULT NULL,
      PRIMARY KEY (id)
      ) ENGINE=InnoDB DEFAULT CHARSET=utf8;
      
      /*Data for the table user */
      #插入数据
      insert  into user(id,name,age) values  ('1','xiaokai',20);
      #修改数据
      update user set age= 30 where id='1';
      #删除数据
      delete from user where id='1';
  • 5、启动kafka的自带控制台消费者

    • 启动Kafka消费者监听maxwell主题
    kafka-console-consumer.sh --topic maxwell --bootstrap-server node01:9092,node02:9092,node03:9092 --from-beginning 
    • 等待一段事件,观察maxwell主题是否有消息发送过来
    {"database":"test_db","table":"user","type":"insert","ts":1621244407,"xid":985,"commit":true,"data":{"id":"1","name":"xiaokai","age":20}}
    
    {"database":"test_db","table":"user","type":"update","ts":1621244413,"xid":999,"commit":true,"data":{"id":"1","name":"xiaokai","age":30},"old":{"age":20}}
    
    {"database":"test_db","table":"user","type":"delete","ts":1621244419,"xid":1013,"commit":true,"data":{"id":"1","name":"xiaokai","age":30}}
    • json数据字段说明

    • database

      • 数据库名称
    • table

      • 表名称
    • type

      • 操作类型
      • 包括 insert/update/delete 等
    • ts

      • 操作时间戳
    • xid

      • 事务id
    • commit

      • 同一个xid代表同一个事务,事务的最后一条语句会有commit
    • data

      • 最新的数据,修改后的数据
    • old

      • 旧数据,修改前的数据

Views: 29

Flume综合案例之拦截器

Flume综合案例之静态拦截器使用

1. 案例场景

  • A、B两台日志服务机器实时生产日志主要类型为access.log、nginx.log、web.log

  • 现在需要把A、B 机器中的access.log、nginx.log、web.log 采集汇总到C机器上然后统一收集到hdfs中。

  • 但是在hdfs中要求的目录为:

/source/logs/access/20200210/**
/source/logs/nginx/20200210/**
/source/logs/web/20200210/**

2. 场景分析

image-20200522145357477

3. 数据流程处理分析

image-20200522145429320

4. 实现

  • 服务器A对应的IP为 192.168.52.100

  • 服务器B对应的IP为 192.168.52.110

  • 服务器C对应的IP为 192.168.52.120

采集端配置文件开发
  • node01与node02服务器开发flume的配置文件
cd /kkb/install/apache-flume-1.9.0-bin/conf/
vim exec_source_avro_sink.conf
  • 内容如下
# Name the components on this agent
a1.sources = r1 r2 r3
a1.sinks = k1
a1.channels = c1

# set source
a1.sources.r1.type = exec
a1.sources.r1.command = tail -F /kkb/install/taillogs/access.log
a1.sources.r1.interceptors = i1
a1.sources.r1.interceptors.i1.type = static
##  static拦截器的功能就是往采集到的数据的header中插入自定义的key-value对;与node03上的agent的sink中的type相呼应
a1.sources.r1.interceptors.i1.key = type
a1.sources.r1.interceptors.i1.value = access

a1.sources.r2.type = exec
a1.sources.r2.command = tail -F /kkb/install/taillogs/nginx.log
a1.sources.r2.interceptors = i2
a1.sources.r2.interceptors.i2.type = static
a1.sources.r2.interceptors.i2.key = type
a1.sources.r2.interceptors.i2.value = nginx

a1.sources.r3.type = exec
a1.sources.r3.command = tail -F /kkb/install/taillogs/web.log
a1.sources.r3.interceptors = i3
a1.sources.r3.interceptors.i3.type = static
a1.sources.r3.interceptors.i3.key = type
a1.sources.r3.interceptors.i3.value = web

# set sink
a1.sinks.k1.type = avro
a1.sinks.k1.hostname = node03
a1.sinks.k1.port = 41415

# set channel
a1.channels.c1.type = memory
a1.channels.c1.capacity = 20000
a1.channels.c1.transactionCapacity = 10000

# Bind the source and sink to the channel
a1.sources.r1.channels = c1
a1.sources.r2.channels = c1
a1.sources.r3.channels = c1
a1.sinks.k1.channel = c1

static interceptor 静态拦截器

服务端配置文件开发
  • 在node03上面开发flume配置文件
cd /kkb/install/apache-flume-1.9.0-bin/conf/
vim avro_source_hdfs_sink.conf
  • 内容如下
# Name the components on this agent
a1.sources = r1
a1.sinks = k1
a1.channels = c1

#定义source
a1.sources.r1.type = avro
a1.sources.r1.bind = node03
a1.sources.r1.port =41415

#定义channels
a1.channels.c1.type = memory
a1.channels.c1.capacity = 20000
a1.channels.c1.transactionCapacity = 10000

#定义sink
a1.sinks.k1.type = hdfs
a1.sinks.k1.hdfs.path=hdfs://node01:8020/source/logs/%{type}/%Y%m%d
a1.sinks.k1.hdfs.filePrefix =events
a1.sinks.k1.hdfs.fileType = DataStream
a1.sinks.k1.hdfs.writeFormat = Text
#时间类型
a1.sinks.k1.hdfs.useLocalTimeStamp = true
#生成的文件不按条数生成
a1.sinks.k1.hdfs.rollCount = 0
#生成的文件按时间生成
a1.sinks.k1.hdfs.rollInterval = 30
#生成的文件按大小生成
a1.sinks.k1.hdfs.rollSize  = 10485760
#批量写入hdfs的个数
a1.sinks.k1.hdfs.batchSize = 10000
#flume操作hdfs的线程数(包括新建,写入等)
a1.sinks.k1.hdfs.threadsPoolSize=10
#操作hdfs超时时间
a1.sinks.k1.hdfs.callTimeout=30000

#组装source、channel、sink
a1.sources.r1.channels = c1
a1.sinks.k1.channel = c1
采集端文件生成脚本
  • 在node01与node02上面开发shell脚本,模拟数据生成
cd /kkb/install/shells
vim server.sh
  • 内容如下
#!/bin/bash
while true
do  
 date >> /kkb/install/taillogs/access.log; 
 date >> /kkb/install/taillogs/web.log;
 date >> /kkb/install/taillogs/nginx.log;
  sleep 0.5;
done
  • node01、node02给脚本添加可执行权限
chmod u+x server.sh
顺序启动服务
  • node03启动flume实现数据收集
cd /kkb/install/apache-flume-1.9.0-bin/
bin/flume-ng agent -c conf -f conf/avro_source_hdfs_sink.conf -name a1 -Dflume.root.logger=DEBUG,console
  • node01与node02启动flume实现数据监控
cd /kkb/install/apache-flume-1.9.0-bin/
bin/flume-ng agent -c conf -f conf/exec_source_avro_sink.conf -name a1 -Dflume.root.logger=DEBUG,console
  • node01与node02启动生成文件脚本
cd /kkb/install/shells
sh server.sh
  • 查看hdfs目录/source/logs

Flume综合案例之自定义拦截器使用

案例需求:

  • 在数据采集之后,通过Flume的拦截器,实现将无效的JSON格式的消息过滤掉。

实现步骤

第一步:创建maven java工程,导入jar包
<dependencies>
    <dependency>
        <groupId>org.apache.flume</groupId>
        <artifactId>flume-ng-core</artifactId>
        <version>1.9.0</version>
    </dependency>
</dependencies>
<build>
    <plugins>
        <plugin>
            <groupId>org.apache.maven.plugins</groupId>
            <artifactId>maven-compiler-plugin</artifactId>
            <version>3.0</version>
            <configuration>
                <source>1.8</source>
                <target>1.8</target>
                <encoding>UTF-8</encoding>
                <!--    <verbal>true</verbal>-->
            </configuration>
        </plugin>
        <plugin>
            <groupId>org.apache.maven.plugins</groupId>
            <artifactId>maven-shade-plugin</artifactId>
            <version>3.1.1</version>
            <executions>
                <execution>
                    <phase>package</phase>
                    <goals>
                        <goal>shade</goal>
                    </goals>
                    <configuration>
                        <filters>
                            <filter>
                                <artifact>*:*</artifact>
                                <excludes>
                                    <exclude>META-INF/*.SF</exclude>
                                    <exclude>META-INF/*.DSA</exclude>
                                    <exclude>META-INF/*.RSA</exclude>
                                </excludes>
                            </filter>
                        </filters>
                        <transformers>
                            <transformer implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer">
                                <mainClass></mainClass>
                            </transformer>
                        </transformers>
                    </configuration>
                </execution>
            </executions>
        </plugin>
    </plugins>
</build>
第二步:自定义flume的拦截器

新建package com.niit.flume.interceptor

新建类及内部类,分别是JsonInterceptorMyBuilder

package com.niit.flume.interceptor;

import org.apache.commons.codec.Charsets;
import org.apache.commons.lang.StringUtils;
import org.apache.flume.Context;
import org.apache.flume.Event;
import org.apache.flume.interceptor.Interceptor;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.util.ArrayList;
import java.util.List;
import java.util.Map;

/**
 * @Author: deLucia
 * @Date: 2021/6/2
 * @Version: 1.0
 * @Description:
 */
public class JsonInterceptor implements Interceptor {
    private static final Logger LOG = LoggerFactory.getLogger(JsonInterceptor.class);
    private String tag = "";

    public JsonInterceptor(String tag) {
        this.tag = tag;
    }

    @Override
    public void initialize() {

    }

    /**
     * 过滤JSON格式之外的数据 {“key":value, {...}}
     * @param event
     * @return
     */
    @Override
    public Event intercept(Event event) {
        String line = new String(event.getBody(), Charsets.UTF_8);
        line = line.trim();
        if (StringUtils.isBlank(line)) {
            return null;
        }

        // {..}
        if (line.startsWith("{") && line.endsWith("}")) {

            Map<String, String> headers = event.getHeaders();
            headers.put("tag", tag);
            return event;

        } else {
            LOG.warn("Not Valid JSON: " + line);
            return null;
        }
    }

    @Override
    public List<Event> intercept(List<Event> list) {
        List<Event> out = new ArrayList<Event>();
        for (Event event : list) {
            Event outEvent = intercept(event);
            if (outEvent != null) {
                out.add(outEvent);
            }
        }
        return out;
    }

    @Override
    public void close() {

    }
    /**
     * 相当于自定义Interceptor的工厂类
     * 在flume采集配置文件中通过指定该Builder来创建Interceptor对象
     *
     * @author
     */
    public static class MyBuilder implements Interceptor.Builder {

        private String tag;

        @Override
        public void configure(Context context) {
            //从flume的配置文件中获得拦截器的“tag”属性值
            this.tag = context.getString("tag", "").trim();
        }

        /*
         * @see org.apache.flume.interceptor.Interceptor.Builder#build()
         */
        @Override
        public JsonInterceptor build() {
            return new JsonInterceptor(tag);
        }
    }
}
第三步:打包上传服务器
  • 将我们的拦截器打成jar包放到主机名为hadoop100的服务器上的flume安装目录下的lib目录下
第四步:开发flume的配置文件

开发flume的配置文件

cd /opt/pkg/flume/conf/
vim nc-interceptor-logger.conf
  • 内容如下

记得将下边的i1.type根据自己的实际情况进行替换

#声明三种组件
a1.sources = r1
a1.channels = c1
a1.sinks = k1

#定义source信息
a1.sources.r1.type=netcat
a1.sources.r1.bind=localhost
a1.sources.r1.port=8888
#定义source的拦截器
a1.sources.r1.interceptors = i1
#根据自定义的拦截器,相应修改全类名及内部类名
a1.sources.r1.interceptors.i1.type =com.niit.flume.interceptor.JsonInterceptor$MyBuilder
## tag的内容将被设置到消息头中
a1.sources.r1.interceptors.i1.tag = json_event

#定义sink信息
a1.sinks.k1.type=logger

#定义channel信息
a1.channels.c1.type=memory

#绑定在一起
a1.sources.r1.channels=c1
a1.sinks.k1.channel=c1
第五步:运行FLume Agent

同时开启INFO日志级别

[hadoop@hadoop100 conf]$ bin/flume-ng agent -c conf/ -f conf/nc_interceptor_logger.conf -n a1 -Dflume.root.logger=INFO,console
第六步:简单测试

通过netcat客户端发送数据

[hadoop@hadoop100 ~]$ nc localhost 8888
asd
dsa
{"OK
OK                                                                             {
OK
OK
{"key":123}
OK

查看flume的agent的日志信息如下

2021-05-26 14:40:12,233 (lifecycleSupervisor-1-4) [INFO - org.apache.flume.source.NetcatSource.start(NetcatSource.java:155)] Source starting
2021-05-26 14:40:17,326 (lifecycleSupervisor-1-4) [INFO - org.apache.flume.source.NetcatSource.start(NetcatSource.java:166)] Created serverSocket:sun.nio.ch.ServerSocketChannelImpl[/127.0.0.1:8888]
2021-05-26 14:40:20,952 (netcat-handler-0) [WARN - com.niit.flume.interceptor.JsonInterceptor.intercept(JsonInterceptor.java:48)] Invalid json format event !
2021-05-26 14:40:20,953 (netcat-handler-0) [WARN - com.niit.flume.interceptor.JsonInterceptor.intercept(JsonInterceptor.java:48)] Invalid json format event !
2021-05-26 14:40:24,109 (netcat-handler-0) [WARN - com.niit.flume.interceptor.JsonInterceptor.intercept(JsonInterceptor.java:48)] Invalid json format event !
2021-05-26 14:40:24,909 (netcat-handler-0) [WARN - com.niit.flume.interceptor.JsonInterceptor.intercept(JsonInterceptor.java:48)] Invalid json format event !
2021-05-26 14:40:36,895 (SinkRunner-PollingRunner-DefaultSinkProcessor) [INFO - org.apache.flume.sink.LoggerSink.process(LoggerSink.java:95)] Event: { headers:{tag=json_event} body: 7B 22 6B 65 79 22 3A 31 32 33 7D                {"key":123} }

Views: 124

大数据日志分析项目需求

目标:电商网站+电商网站后台管理系统+大数据分析+数据可视化思路:

按照数据的采集,数据的存储,数据分析处理,数据可视化

逻辑图:

项目要求   

  • 题材不限,但需先经过老师认可
  • 各组组长记录组员项目进度,每周提交给老师
  • 数据源至少来自两处,可以是日志、关系型数据库、以及爬虫的数据
  • 日志允许通过代码生成(或者ab压测工具来生成)
  • 前后台管理页面不能和老师的一样   
  • 讲述清楚nginx、tomcat、flume、sqoop、hadoop各自作用   
  • 重点突出数据分析部分,mapreduce、hive、storm至少使用2种,完成6个不同的分析,hbase、kafka选择使用
  • 日志收集使用flume, 数据导入导出使用sqoop
  • 数据可视化可以使用echarts、或其他类型前端框架、另外superset,datav等也允许使用
  • 能够熟练使用各类脚本及命令
  • 最后一周(17周结束前)完成所有项目演讲,每人1-2分钟时间
  • 大数据采集、存储、分析处理、可视化的过程应该是连续的
  • 允许使用任务调度框架编排执行定时任务
  • 展示的时候所有模块都需要完全部署到虚拟机,不允许本地运行项目

七、开源的爬虫项目   

网站:http://www.geccocrawler.com/tag/sysc/    GitHub:https://github.com/xtuhcy/gecco

Views: 141

分布式调度系统-Apache DolphinScheduler(集群部署)

一、课前准备

  • Hadoop-3.1.2集群
  • MySQL-5.7
  • zookeeper-3.6.2集群
  • Hive-3.1.2
  • spark-2.3.3

二、课堂目标

  • 熟练使用DolphinScheduler调度系统

三、知识要点

1、DolphinScheduler简介

2、DolphinScheduler的特性

2.1 高可靠性
  • 去中心化多Master和Worker,自身支持HA功能,采用任务队列来避免过载,不会造成机器卡死

2.2.简单易用

  • DAG监控界面,所有流程定义都是可视化,通过拖拽任务制定DAG
  • 通过API方式与第三方系统对接,一键部署。

2.3.丰富的使用场景

  • 支持暂停、恢复操作,支持多租户,更好的应对大数据的使用场景,支持更多的任务类型,如hive,mr,spark,python

2.4.高扩展性

  • 支持自定义任务类型,调度器使用分布式调度,调度能力随集群线性增长,Master和Worker支持动态上下线

3、DolphinScheduler的架构介绍

3.1 系统架构设计

https://dolphinscheduler.apache.org/zh-cn/blog/architecture-design.html

3.2 DS-1.3改进及新特性

  • 数据库减压,减少极端情况下的可能造成的调度延时

  • Worker去DB、职责更单一

  • Master和Worker直接通信,降低 延时

  • Master多种策略分发任务(有三种方式选择Worker节点:随机、循环、CPU和 内存的线性加权负载平衡 )

  • 资源中心支持多目录

  • 任务类型新增Datax、 Sqoop、条件分支

  • DAG一键格式化

  • 批量导出和导入工作流

  • 工作流复制

3.3 DS-1.3系统架构图

image-20210326163738309

4、前置环境准备

说明:

安装DolphinScheduler(以下简称ds)前,建议跟文档下边说明的环境保持统一

否则安装及使用ds的过程中,可能会出现位置错误需要自己解决

4.1 安装JDK-1.8
  • 此文档以3节点在/kkb/install都安装了jdk1.8.0_141为例进行演示
  • 安装包jdk-8u141-linux-x64.tar.gz
4.2 安装Hadoop-3.1.4集群
4.3 安装Zookeeper-3.6.2集群
4.4 安装Mysql-5.7
4.5 安装hive-3.1.2
  • node02、node03节点安装了hive-3.1.2
  • 若没有如此安装,参考资料《Hive安装部署》进行安装
4.6 安装Spark-2.3.3
  • dolphinscheduler中会演示调度spark程序,先演示基本的用法

  • 此文档以3节点在/kkb/install都安装了spark-2.3.3为例进行演示

  • 若没有如此安装,参考资料《spark安装部署.md》进行安装

5、安装部署

官方安装指导:https://dolphinscheduler.apache.org/zh-cn/docs/1.3.4/user_doc/quick-start.html

5.1节点规划
机器 服务 端口 group
node01 master、api、logger 8787(master)、8888(api)
node02 master、alert、worker、logger 8787(master)、7878(worker) hadoop
node03 worker、logger 7878(worker) hadoop

hadoop 组配置后master分发任务才能根据cpu和内存的负载选择具体哪个worker执行任务

5.2 准备工作
1、创建目录
  • 确保三个节点都有目录/kkb/soft/kkb/install,且所属用户及用户组如下
[hadoop@node01 ~]$ ll /kkb/ 
总用量 0
drwxr-xr-x. 2 hadoop hadoop 6 4月  13 14:21 install
drwxr-xr-x. 2 hadoop hadoop 6 4月  13 14:22 soft

image-20210413142651607

  • 若没有这些目录,那么如下创建;==3个节点==都运行如下命令
sudo mkdir -p /kkb/install
sudo mkdir -p /kkb/soft
sudo chown -R hadoop:hadoop /kkb/
ll /kkb/
  • 确保目录所属变成如下样子

image-20210413142651607

2、确保已安装zookeeper集群
  • 启动zookeeper集群
    • 确保三节点上已经安装了zk;
    • zk版本要求:ZooKeeper (3.4.6+)
    • 若没有安装,请先安装再往下继续
3、启动HDFS
  • 因为ds的资源存储在HDFS上

  • 所以,node01上运行start-dfs.sh启动hdfs

5.3 开始安装

第一步:安装包下载,在node01上执行

[hadoop@node01 ~]$ cd /kkb/soft/
[hadoop@node01 soft]$ wget https://mirrors.tuna.tsinghua.edu.cn/apache/dolphinscheduler/1.3.5/apache-dolphinscheduler-incubating-1.3.5-dolphinscheduler-bin.tar.gz

第二步:解压压缩包

[hadoop@node01 soft]$ tar -xzvf apache-dolphinscheduler-incubating-1.3.5-dolphinscheduler-bin.tar.gz -C /kkb/install/

第三步:重命名

[hadoop@node01 soft]$ cd /kkb/install/
[hadoop@node01 install]$ mv apache-dolphinscheduler-incubating-1.3.5-dolphinscheduler-bin/ dolphinscheduler-1.3.5
[hadoop@node01 install]$ ll
总用量 0
drwxrwxr-x. 9 hadoop hadoop 156 4月  13 14:35 dolphinscheduler-1.3.5

第四步:建库建表

此处以node03安装了mysql为例

node01上

[hadoop@node01 install]$ scp /kkb/install/dolphinscheduler-1.3.5/sql/dolphinscheduler_mysql.sql node03:/kkb/soft/

node03上,进入MySQL命令行执行

mysql -uroot -p

set global validate_password_policy=LOW;
set global validate_password_length=6;
CREATE DATABASE IF NOT EXISTS dolphinscheduler DEFAULT CHARSET utf8 DEFAULT COLLATE utf8_general_ci;
GRANT ALL PRIVILEGES ON dolphinscheduler.* TO 'root'@'%' IDENTIFIED BY '123456';
GRANT ALL PRIVILEGES ON dolphinscheduler.* TO 'root'@'localhost' IDENTIFIED BY '123456';
flush privileges;
use dolphinscheduler;
source /kkb/soft/dolphinscheduler_mysql.sql;
mysql> show tables;
+--------------------------------+
| Tables_in_dolphinscheduler     |
+--------------------------------+
| QRTZ_BLOB_TRIGGERS             |
| QRTZ_CALENDARS                 |
| QRTZ_CRON_TRIGGERS             |
| QRTZ_FIRED_TRIGGERS            |
| QRTZ_JOB_DETAILS               |
| QRTZ_LOCKS                     |
| QRTZ_PAUSED_TRIGGER_GRPS       |
| QRTZ_SCHEDULER_STATE           |
| QRTZ_SIMPLE_TRIGGERS           |
| QRTZ_SIMPROP_TRIGGERS          |
| QRTZ_TRIGGERS                  |
| t_ds_access_token              |
| t_ds_alert                     |
| t_ds_alertgroup                |
| t_ds_command                   |
| t_ds_datasource                |
| t_ds_error_command             |
| t_ds_process_definition        |
| t_ds_process_instance          |
| t_ds_project                   |
| t_ds_queue                     |
| t_ds_relation_datasource_user  |
| t_ds_relation_process_instance |
| t_ds_relation_project_user     |
| t_ds_relation_resources_user   |
| t_ds_relation_udfs_user        |
| t_ds_relation_user_alertgroup  |
| t_ds_resources                 |
| t_ds_schedules                 |
| t_ds_session                   |
| t_ds_task_instance             |
| t_ds_tenant                    |
| t_ds_udfs                      |
| t_ds_user                      |
| t_ds_version                   |
+--------------------------------+
35 rows in set (0.00 sec)

第五步:修改配置文件

[hadoop@node01 ~]$ cd /kkb/install/dolphinscheduler-1.3.5/conf/

1、alert.properties #配置告警邮箱相关信息

  • 此处已126邮箱为例
  • 登录自己的126邮箱

image-20210415113539557

image-20210415113712981

image-20210415113809874

image-20210415113906134

image-20210415114129004

  • 开始配置文件
[hadoop@node01 conf]$ vim alert.properties
#alert type is EMAIL/SMS
alert.type=EMAIL

# mail server configuration
mail.protocol=SMTP
mail.server.host=smtp.126.com
mail.server.port=25
mail.sender=youhy964@126.com
mail.user=youhy964@126.com
mail.passwd=WBNPUGCNZMQQYBUT
# TLS
mail.smtp.starttls.enable=true
# SSL
mail.smtp.ssl.enable=false
mail.smtp.ssl.trust=smtp.126.com

2、application-api.properties

  • 修改web ui端口号
[hadoop@node01 conf]$ vim application-api.properties 
# server port
server.port=8888

image-20210413153616572

3、common.properties

# 修改如下3个属性的值
[hadoop@node01 conf]$ vim common.properties
resource.storage.type=HDFS
fs.defaultFS=hdfs://node01:8020
yarn.application.status.address=http://node01:8088/ws/v1/cluster/apps/%s

4、datasource.properties

假设mysql安装在node03节点

# 修改如下几个属性的值
[hadoop@node01 conf]$ vim datasource.properties

spring.datasource.driver-class-name=com.mysql.jdbc.Driver
spring.datasource.url=jdbc:mysql://node03:3306/dolphinscheduler?useUnicode=true&characterEncoding=UTF-8&allowMultiQueries=true
spring.datasource.username=root
spring.datasource.password=123456

5、master.properties

[hadoop@node01 conf]$ vim master.properties

master.listen.port=8787 

6、worker.properties

[hadoop@node01 conf]$ vim worker.properties 
worker.listen.port=7878
worker.groups=hadoop

image-20210413155102581

7、zookeeper.properties

[hadoop@node01 conf]$ vim zookeeper.properties

zookeeper.quorum=node01:2181,node02:2181,node03:2181

8、env/dolphinscheduler_env.sh

# 修改如下属性,根据自己的实际情况,配置属性值
[hadoop@node01 conf]$ vim env/dolphinscheduler_env.sh

export HADOOP_HOME=/kkb/install/hadoop-3.1.4
export HADOOP_CONF_DIR=/kkb/install/hadoop-3.1.4/etc/hadoop
export SPARK_HOME1=/kkb/install/spark-2.3.3-bin-hadoop2.7
export JAVA_HOME=/kkb/install/jdk1.8.0_141
export HIVE_HOME=/kkb/install/apache-hive-3.1.2

image-20210414114035151

第六步:将hdfs-site.xml、core-site.xml 拷贝至ds的 conf 目录下,同时 把mysql-connector-java-5.1.48-bin.jar驱动包上传到ds的lib目录下

[hadoop@node01 conf]$ cp /kkb/install/hadoop-3.1.4/etc/hadoop/hdfs-site.xml /kkb/install/dolphinscheduler-1.3.5/conf
[hadoop@node01 conf]$ cp /kkb/install/hadoop-3.1.4/etc/hadoop/core-site.xml /kkb/install/dolphinscheduler-1.3.5/conf
# 将mysql-connector-java-5.1.38.jar上传到node01的/kkb/soft目录
[hadoop@node01 soft]$ cd /kkb/soft/
[hadoop@node01 soft]$ cp mysql-connector-java-5.1.38.jar /kkb/install/dolphinscheduler-1.3.5/lib/

第七步:到ds的bin目录下(1.3.5版本,可以跳过此步)

dos2unix dolphinscheduler-daemon.sh   #dos2unix:将DOS格式的文本文件转换成UNIX格式的(DOS/MAC to UNIX text file format converter)
chmod +x dolphinscheduler-daemon.sh

第八步:scp到node02、node03

[hadoop@node01 bin]$ cd /kkb/install/
[hadoop@node01 install]$ scp -r dolphinscheduler-1.3.5/ node02:$PWD
[hadoop@node01 install]$ scp -r dolphinscheduler-1.3.5/ node03:$PWD

部署过程中的问题

建表时出现 Index column size too large. The maximum column size is 767 bytes.(此问题在MySQL-5.6及以前版本会出现)
将建表文件中 CHARSET=utf8mb4 改为 CHARSET=utf8
5.4 调优配置

生产环境上建议,worker.properties里设置的cpu和内存调一下就可以保护worker不至于挂掉,一般别超过cpu核数的2倍,内存留上1~2G(当然如果比较豪,可以预留更多资源),线程数别超过cpu核数的2.5倍。

例如:8c16G机器

worker.exec.threads=20
worker.max.cpuload.avg=16
worker.reserved.memory=1

同理对于master调优配置 8c16G机器
master.properties文件

master.max.cpuload.avg=16
master.reserved.memory=1
5.5 启动

在node01 启动 master、api、logger

[hadoop@node01 ~]$ cd /kkb/install/dolphinscheduler-1.3.5/bin/
[hadoop@node01 bin]$ ./dolphinscheduler-daemon.sh start master-server
[hadoop@node01 bin]$ ./dolphinscheduler-daemon.sh start api-server
[hadoop@node01 bin]$ ./dolphinscheduler-daemon.sh start logger-server

image-20210414141821598

在node02 启动 master、alert、worker、logger

[hadoop@node02 ~]$ cd /kkb/install/dolphinscheduler-1.3.5/bin/
[hadoop@node02 bin]$ ./dolphinscheduler-daemon.sh start master-server
[hadoop@node02 bin]$ ./dolphinscheduler-daemon.sh start alert-server
[hadoop@node02 bin]$ ./dolphinscheduler-daemon.sh start worker-server
[hadoop@node02 bin]$ ./dolphinscheduler-daemon.sh start logger-server

image-20210414142819920

在node03 启动 worker、logger

[hadoop@node03 ~]$ cd /kkb/install/dolphinscheduler-1.3.5/bin/
[hadoop@node03 bin]$ ./dolphinscheduler-daemon.sh start worker-server
[hadoop@node03 bin]$ ./dolphinscheduler-daemon.sh start logger-server

image-20210414143002703

目前3节点已经启动了hadoop集群、zookeeper集群、及ds相关进程;

启动成功之后,会出现如下进程

[hadoop@node01 bin]$ xcall jps
============= node01 jps =============
3680 ApiApplicationServer
2465 DataNode
3761 LoggerServer
3410 JobHistoryServer
2036 QuorumPeerMain
2294 NameNode
3607 MasterServer
2666 SecondaryNameNode
2892 ResourceManager
3037 NodeManager
4351 Jps
============= node02 jps =============
2320 NodeManager
2880 LoggerServer
2179 DataNode
3446 Jps
3287 AlertServer
2616 MasterServer
2798 WorkerServer
2031 QuorumPeerMain
============= node03 jps =============
2160 DataNode
2289 NodeManager
3009 LoggerServer
1996 QuorumPeerMain
2927 WorkerServer
3103 Jps
  • 如何关闭ds集群?
    • 关闭个节点的各角色,就是将上边启动命令中的start换成stop即可
[hadoop@node01 ~]$ cd /kkb/install/dolphinscheduler-1.3.5/bin/
[hadoop@node01 bin]$ ./dolphinscheduler-daemon.sh stop master-server
[hadoop@node01 bin]$ ./dolphinscheduler-daemon.sh stop api-server
[hadoop@node01 bin]$ ./dolphinscheduler-daemon.sh stop logger-server

[hadoop@node02 ~]$ cd /kkb/install/dolphinscheduler-1.3.5/bin/
[hadoop@node02 bin]$ ./dolphinscheduler-daemon.sh stop master-server
[hadoop@node02 bin]$ ./dolphinscheduler-daemon.sh stop alert-server
[hadoop@node02 bin]$ ./dolphinscheduler-daemon.sh stop worker-server
[hadoop@node02 bin]$ ./dolphinscheduler-daemon.sh stop logger-server

[hadoop@node03 ~]$ cd /kkb/install/dolphinscheduler-1.3.5/bin/
[hadoop@node03 bin]$ ./dolphinscheduler-daemon.sh stop worker-server
[hadoop@node03 bin]$ ./dolphinscheduler-daemon.sh stop logger-server

6、快速上手

6.1 web页面

image-20210414152427370

  • 登录后,进入如下界面

image-20210414152550281

6.2 首页
  • 点击上图①,进入“首页”,能够看到任务状态、流程状态

image-20210414152702608

  • 任务状态:
    • 某个任务的各种状态;
    • 如上图:提交成功、正在运行、准备暂停、暂停......
    • 通过各种状态,可以看到集群的繁忙程度,尤其是“等待线程”的数量
  • 流程状态:
    • 流程由一个或多个任务组成
6.3 项目管理

image-20210415144646143

  • 点击导航栏“项目管理”
  • 用来创建项目及列出已有项目
6.4 资源管理

image-20210415144719980

  • 可以创建文件夹
  • 然后在文件夹中创建脚本
  • 以后的项目里的流程中的任务可以用到这些脚本
  • 函数管理:主要是创建hive的UDF函数

image-20210415144912359

6.5 数据源中心

image-20210415145058616

  • 可以贴sql操作
6.6 监控中心
  • 可以查看master、worker、zookeeper、db等统计信息
  • master
    • 集群负载小于10比较正常,如果大于10,就比较繁忙了,意味着机器的性能比较弱了

image-20210415145528740

  • worker

image-20210415151542512

  • zookeeper

image-20210415151757429

  • DB

image-20210415151820549

6.7 安全中心

image-20210415151932577

  • 管理员常用的配置页面

7、演示

7.1 安全中心

####### 1、令牌管理

  • 安全中心 - 令牌管理
  • 令牌管理:
    • 在工作中一般用不到
    • java程序跟ds调度器进行集成,用rest api进行交互的时候,会用到令牌

image-20210415152538936

image-20210415152733911

  • 把token给其他项目组,其他项目组进行源代码的二次开发的时候,进行插件集成的时候,会用到token值
  • 工作中正常使用ds调度时,一般用不到
2、队列管理
  • 指的是yarn调度器中的队列或cdh的动态资源池
  • 提交应用到yarn集群时,可以指定将应用提交到指定的队列中,此队列占据集群的一定量的内存、cpu资源,共此应用使用;具体的可以回顾yarn中的调度器相关的知识

image-20210415154518353

  • 假设yarn的root队列下还有个队列叫dw
  • 那么在ds的ui界面创建队列

image-20210415154826352

  • 效果如下

image-20210415154850916

3、Worker分组管理
  • 默认只有default一个分组

image-20210415155200106

  • 如果将来ds集群有一些机器是高配,一些是低配的,那么可以进行分组,比如提交时,指定提交到高配的机器上执行
4、告警组管理
  • 默认如下

image-20210415155408815

  • 创建告警组

image-20210415155542189

  • 管理告警组关联的管理用户

image-20210415155748320

  • 随后创建好一个ds用户后,再在上图中关联“管理用户” todo
5、租户管理
  • 如果要创建ds的用户,会用到租户
  • 所以,此处,我们先创建租户
  • 租户:
    • ds调度执行时,会通过su - hadoop类似的用法,切换到linux的某用户,进行执行任务

image-20210415162157184

  • 效果如下

image-20210415162245361

6、用户管理

image-20210415162725052

  • 效果如下

image-20210415162822186

  • 接下来去“告警组管理”中,给刚创建的告警组“kkb warning”进行管理用户授权,如下操作
  • 告警发生后,发送给谁?此处我们发送给kkbds用户的邮箱

image-20210415163007041

7.2 项目管理
1、创建项目

image-20210415164314961

image-20210415164342212

image-20210415164406304

2、创建工作流
  • 进入项目页面

image-20210415164438803

image-20210415164609657

image-20210415164628972

image-20210415164824273

  • 拖动shell图标到右侧区域,添加一个shell类型的task,并在下图中设置相关的参数

image-20210415165019560

  • 复制一个task

image-20210415165124319

  • 给新task做参数设置

image-20210415165603515

  • 串联两个task,并保存

image-20210415165911187

  • 弹出下图

image-20210415170044076

  • 出现下图效果

image-20210415170426952

  • 流程按钮说明

image-20210415171529417

  • 点击“运行”按钮

image-20210415172118331

3、工作流实例

image-20210415172743087

  • 下图中有两个task,如果task还正在运行,可以点击右上角“刷新”按钮,刷新其状态

image-20210415172845277

image-20210415172949561

  • 效果如下

image-20210415173017674

image-20210415173619809

  • 查看日志

image-20210415173750442

  • 查看历史

image-20210415173849464

4、任务实例

image-20210415174529225

  • 看到firstflow中的两个任务实例都执行成功
  • 去邮箱确认是否收到邮件

image-20210415174655301

7.3 调度hive
1、hive环境准备

环境说明:

  • 此处以==node03安装mysql-5.7、hive-3.1.2为例==
  • ==node02安装了hive-3.12.为例==
  • 若没有安装,请参考资料《Hive安装部署.md》
  • ==启动Hiveserver2==
[hadoop@node03 ~]$ cd /kkb/install/apache-hive-3.1.2/bin
[hadoop@node03 bin]$ ./hiveserver2
  • 若启动hiveserver2时,报错如下
Exception in thread "main" java.lang.NoSuchMethodError: com.google.common.base.Preconditions.checkArgument(ZLjava/lang/String;Ljava/lang/Object;)V
  • 原因是hadoop、hive中的guava包冲突
  • 解决方案:将hive的此包删除,hadoop的guava包复制一个到hive的lib目录
[hadoop@node03 ~]$ cd /kkb/install/apache-hive-3.1.2/lib/
[hadoop@node03 lib]$ rm -rf guava-19.0.jar
[hadoop@node03 lib]$ cp /kkb/install/hadoop-3.1.4/share/hadoop/common/lib/guava-27.0-jre.jar ./
[hadoop@node03 lib]$
2、方式一:贴hive sql

1、创建hive数据源

  • 如下图操作

image-20210416141935180

  • 点击上图右下角“提交”按钮

  • 效果如下

image-20210416142800187

  • 去“项目管理”,打开前边创建的项目"kkbpro"

image-20210416142929579

image-20210416143057547

  • 进入下图,将task01、task02禁用掉
  • 以task01操作为例:

image-20210416143348069

先在hive的default中创建表,并插入数据,用于一会ds中测试hive数据源操作

create table stu(id int, name string);
insert into stu values(1, "zs");
insert into stu values(2, "ww");
查询结果
hive (default)> select * from stu;
OK
stu.id    stu.name
1 zs
2 ww

image-20210416144535396

image-20210416144605198

image-20210416144640821

image-20210416144816537

image-20210416144911579

  • 查看工作流实例

image-20210416145039900

image-20210416145208006

image-20210416145221814

  • 查看邮件

image-20210416145519318

3、方式二:调度hive脚本

1、hive脚本操作

  • 上边讲了hive数据源中,贴sql的操作
  • 接下来看看如果调用hive脚本
  • 先给ds的admin用户添加租户

image-20210416150450204

image-20210416150511036

  • 资源中心相关操作

image-20210416150613096

image-20210416150639422

image-20210416150733385

image-20210416150951832

创建脚本

image-20210416151655692

文件内容如下

#!/bin/bash

SQL="SELECT * FROM default.stu"

echo "operate with hive cli---->"

hive -e "${SQL}"

echo "operate with beeline---->"
/kkb/install/apache-hive-3.1.2/bin/beeline -u "jdbc:hive2://node03:10000/default" hive -e "${SQL}"

2、调用脚本

  • 项目管理 -》kkbpro项目

image-20210416151841289

  • 将firstflow工作流下线,并重新编辑

image-20210416151915462

image-20210416152027166

  • 让后工作流中添加shell脚本类型的task,task直行上边创建的hive脚本文件

image-20210416152331044

  • 保存工作流
  • 重新上线工作流,然后运行(这些操作之前都已经演示过,不再重复给出截图)

image-20210416152519254

  • 查看工作流实例

image-20210416152833328

  • 查看hive-shell任务的日志

image-20210416152955286

7.4 调度spark
1、spark环境准备
  • 若spark还未学习,可以先看此部分的讲解视频,不进行实操
  • 因为ds可以调度spark,所以此文档先将内容给出

环境说明:

所有的ds worker节点都安装了spark

  • 先启动spark集群
  • node03节点,启动hive的metastore服务
hive --service metastore
  • 启动spark thriftserver服务
[hadoop@node03 sbin]$ cd /kkb/install/spark-2.3.3-bin-hadoop2.7/sbin
[hadoop@node03 sbin]$ ./start-thriftserver.sh \
--hiveconf hive.server2.thrift.port=10000 \
--master yarn-client \
--driver-cores 1 \
--driver-memory 1G \
--executor-cores 1 \
--executor-memory 1G \
--num-executors 2

​ 以yarn-client的方式运行spark thriftserver

​ yarn界面

image-20210416162536120

2、方式一:直接贴sql
  • 数据源中心-创建数据源

image-20210416165625761

image-20210416165733051

  • 进入项目管理 -》 项目kkbpro -》工作流定义 -》下线firstflow工作流 -》编辑firstflow工作流 -》hive-shell 任务停用

image-20210416170041413

  • 然后创建一个新的sql任务,并配置

image-20210416170406671

  • 保存工作流

image-20210416170509205

  • 工作流重新上线,并运行

image-20210416170624387

  • 查看此工作流的工作流实例中spark-datasource任务的日志

image-20210416170858767

  • 查看通知邮件

image-20210416170941093

3、方式二:执行spark脚本
  • 类似于hive脚本操作
  • 通过spark beeline操作脚本
  • 资源中心-创建文件夹

image-20210416175522890

image-20210416175546371

image-20210416175611504

image-20210416175723034

  • 脚本内容如下
#!/bin/bash

SQL="SELECT * FROM default.stu"

echo "operate with spark beeline------------->->->"
/kkb/install/spark-2.3.3-bin-hadoop2.7/bin/beeline \
-u "jdbc:hive2://node03:10000/default" hive -e "${SQL}"
  • 效果如下

image-20210416175817624

  • 如下类似的操作,之前已经做过截图,所以部分操作截图省略

  • 进入项目管理 -》 项目kkbpro -》工作流定义 -》下线firstflow工作流 -》编辑firstflow工作流 -》spark-datasource任务停用

image-20210416180205803

  • 工作流中添加SHELL类型的任务

image-20210416180410253

  • 保存工作流
  • 重新上线firstflow工作流,并运行

image-20210416180530791

  • 查看日志及邮件

image-20210416180637367

image-20210416180738193

Views: 47