imay closed pull request #263: Tidy up the docs and gensrc directory URL: https://github.com/apache/incubator-doris/pull/263
This is a PR merged from a forked repository. As GitHub hides the original diff on merge, it is displayed below for the sake of provenance: As this is a foreign pull request (from a fork), the diff is supplied below (as it won't show otherwise due to GitHub magic): diff --git a/docs/design/Palo_privilege.md b/docs/design/Palo_privilege.md deleted file mode 100644 index f3072906..00000000 --- a/docs/design/Palo_privilege.md +++ /dev/null @@ -1,279 +0,0 @@ -# Palo权限管理 - -## 问题 -1. 当前对Palo中数据的访问权限控制仅到 DB 级别,无法更细粒度的控制到表甚至列级别的访问权限。无法满足部分用户对于数据仓库权限管理的需求。 -2. 当前白名单机制和账户管理模块各自独立,而实际这些都数据权限管理模块,应该整合。 -3. 用户角色无法满足云上部署需求。Root用户拥有所有权限,既可以操作集群启停,又可以访问用户数据,无法满足在云上同时进行集群管理和屏蔽用户数据访问权限的需求。 - -## 名词解释 - -1. 权限管理 - - 权限管理主要负责控制,对于满足特定标识的行为主体,对数据仓库中某一个(或一类)特定元素,可以执行哪些操作。比如来自 127.0.0.1 的 Root 用户(行为主体),可以对 Backend 节(元素)点进行添加或删除操作。 - -2. 账户管理 - - 账户管理主要负责控制,允许哪些满足特定标识的行为主体对数据仓库进行连接和操作。包括密码检查等。 - -3. 权限(Privilege) - - 数据库中定义个多种权限,如Select、Insert 等,用于描述对于对应操作的权限。 - -4. 角色(Role) - - 角色是一组权限的集合 - -5. 用户(User) - - 用户以 Who@Where 的方式表示。如 [email protected] 表示来自 127.0.0.1,名为root的用户。用户可以被赋予权限,已可以归属于某个角色。 - -## 详细设计 - -### 操作 -Palo 中的操作主要分为以下几类: - -1. 节点管理 - - 包括各类节点的添加、删除、下线等。以及对cluster的创建、删除、更改 - -2. 数据库管理 - - 包括对数据库的创建、删除及更改操作。 - -3. 表管理 - - 包括在数据库内,对表的创建、删除及更改操作。 - -4. 读写 - - 对用户数据的读写操作。 - -5. 权限操作 - - 包括创建、删除用户,授予、撤销权限 - -6. 备份恢复操作 - - 对数据库表的备份恢复操作。 - -### 权限 - -#### Node_priv - -节点级别的操作权限 -对应操作:Add/Drop/Decommission Frontend/Backend/Broker - -#### Grant_priv - -进行权限操作以及添加、删除用户的权限。拥有该权限的用户,只能对自身拥有的权限进行授予、撤回等操作。 - -#### Select_priv - -对用户数据的读取权限。 -如果赋予DB级别,则可以读取该DB下的所有数据表。如果赋予Table级别,则仅可以读取该Table中的数据。对于没有 Select\_priv 权限的表,用户无法访问,也不可见。 - -#### Load_priv -因为Palo目前只有导入这一种数据写入方式,所以不再区分Insert、Update、Delete这些细分权限。用户可以对有Load\_priv的表进行数据更改操作。但如果没有对该表的 Select\_priv 权限,则不可以读取该表的数据,但可以查看该表的Schema等信息。 - -#### Alter_priv - -修改DB或者Table的权限。该权限可以修改对应DB或者Table的Schema。同时拥有查看Schema的权限。但是不能进行DB或Table的创建和删除操作。 - -#### Create_priv - -创建DB或者Table的权限。同时拥有查看Schema的权限。 - -#### Drop_priv - -删除DB或者Table的权限。同时拥有查看Schema的权限。 - -### 角色 - -#### Root - -为了和之前的代码保持兼容,这里沿用Root这个名称。实际上,该角色为集群管理员,只拥有 Node\_priv 权限。在Palo集群创建时,默认会创建一个名为 root@'%' 的用户。该角色有且仅有一个用户,不可创建、更改或删除。 - -#### Superuser - -该角色拥有除 Node\_priv 以外的所有权限,意味管理员角色。在Palo集群创建时,默认会创建一个 Admin@'%'的用户。因为Superuser拥有 Grant\_priv,所以可以创建其他的拥有除Node\_priv以外的任意权限的角色,包括创建新的 Superuser。 - -#### Other Roles - -Root 和 Superuser 角色为保留角色。用户可以自定义其他角色。其他角色不可包含 Node\_priv 和 Grant\_priv 权限。 - -## 数据结构 - -参照Mysql的权限表组织方式。我们生成默认的名为mysql数据库,其中包含以下数据表 - -1. user - -| Host | User | Password | Node\_priv | Grant\_priv | Select\_priv | Load\_priv | Alter\_priv | Create\_priv | Drop\_priv | -|---|---|---|---|---|---|---|---|---|---|---| -| % | root |*B371DC178FD00FA65915DBC87B05C7E88E3BE66A| Y | N | N | N | N | N | N | N | - -Host 和 User 标识连接主体,通过 Host、User、Password 来验证连接是否合法。 -该表中的权限问全局权限(Global Priv),全局权限可以作用于任何数据库中的任何表。 - -2. db - -| Host | Db | User | Select\_priv | Load\_priv | Alter\_priv | Create\_priv | Drop\_priv | -|---|---|---|---|---|---|---|---|---|---|---| -| % | example_db | cmy | Y | Y | N | N | N | - -该表用于存储DB级别的权限,通过 Host、DB、User 匹配DB级别的权限。 - -3. tables_priv - -| Host | Db | User | Table_name| Table_priv | Column_priv | -|---|---|---|---|---|---|---|---|---|---|---| -| % | example_db | cmy | example_tbl | Select\_priv/Load\_priv/Alter\_priv/Create\_priv/Drop\_priv | Not support yet | - -## 权限检查规则 - -### 连接 - -通过 User 表进行连接鉴权。User 表按照 Host、User 列排序,取最先匹配到的Row作为授予权限。 - -### 请求 - -连接成功后,用户后续发送的请求需要进行鉴权。首先在User表中查看全局权限,如果没有找到,则再一次查找 db、tables_priv 表。当然如 Node\_priv 这种权限只会出现在 User 表中,则不再继续后面的查找了。 - - -## 语法 - -### 创建用户 - -``` -CREATE USER [IF NOT EXIST] 'user_name'@'host' IDENTIFIED BY 'password' DEFAULT ROLE 'role_name'; -``` - -1. user\_name: 用户名 -2. host: 可以是ip、hostname。允许通配符,如:"%", "192.168.%", "%.baidu.com"(如何解析BNS??) -3. 如果设置了 default role,则直接赋予default role 对应的权限。 - -### 删除用户 - -``` -DROP USER 'user_name'; -``` - -删除对应名称的用户,该用户的所有对应host的权限信息都会被删除。 - - -### 设置密码 -``` -SET PASSWORD FOR 'user_name'@'host' = PASSWORD('xxxx'); -``` - -### 授权 - -``` -GRANT PRIVILEGE 'priv1'[, 'priv2', ...] on `db`.`tbl` to 'user'@'host'[ROLE 'role']; -``` - -``` -GRANT ROLE 'role1' on `db`.`tbl` to 'user'@'host'; -``` - -1. db 和 tbl 可以是通配符,如:\*.\*, db.* 。 Palo不检查执行授权时,db或tbl是否存在。只是更新授权表。 -2. Node\_priv 和 Grant\_priv 不能通过该语句授权。不能授予 all privileges -3. 如果是授权给 Role,且Role不存在,会自动创建这个 Role。如果user不存在,则报错。 -4. 可以授权 ROLE 给某一个user。 - -### 撤权 - -``` -REVOKE 'priv1'[, 'priv2', ...]|[all privileges] on `db`.`tbl` from 'user'@'host'[ROLE 'role']; -``` - -1. db 和 tbl 可以是通配符,Palo会匹配权限表中的所有可匹配项,并撤销对应权限 -2. 该语句中的 host 会进行精确匹配,而不是通配符。 - -### 创建角色 - -``` -CREATE ROLE 'role'; -``` - -1. 创建角色后,可以通过授权语句或撤权语句对该 Role 的权限进行修改。 -2. 默认有 ROOT 和 SUPERUSER 两个 ROLE。这两个ROLE的权限不可修改。 - - -## 最佳实践 - -1. 创建集群后,会自动创建 ROOT@'%' 和 ADMIN@'%' 两个用户。以及 ROOT 和 SUPERUSER 两个角色。 ROOT@'%' 为 ROOT 角色。ADMIN@'%' 为 SUPERUSER 角色。 -2. ROOT 角色仅拥有 Node_priv。SUPERUSER 角色拥有其他所有权限。 -3. ROOT@'%' 主要由于运维人员搭建系统。或者用于云上的部署程序。 -4. SUPERUSER 角色的用户为数据库的实际使用者和管理员。 -5. superuser 用户可以创建各个普通用户。比如给外包人员创建 Alter_priv、Load_priv、Create_priv、Drop_priv 权限,用于数据创建和导入,但没有Select_priv。给用户开通 Select_priv,可以查看数据,但不能对数据进行修改。 -6. 初始时,只有 superuser 可以创建db。或者可以先授权某个普通用户对一个**不存在的DB**的创建权限,之后,该普通用户就可以创建数据库了。(这是一个先有鸡还是先有蛋的问题) -7. 任何对数据库内对象(DB、Table)的创建和删除都不会影响已经存在的权限。如果有必要,必须手动修改权限表对应的条目。 - - -## 排期 -预计6月底完成,本期实现表级别的权限管理。 - -## 遗留问题 -1. 白名单,如何解决bns和dns的问题 - 保留当前通过 add whitelist 的方式添加白名单的机制。通过这个机制添加的白名单,默认都是DNS或BNS。Palo会有后台线程定期解析这些域名,将解析出来的host或ip加入到权限列表中。加入权限表时,取此user在db权限表中,各db最大权限集合。 - - 这种方式,能够兼容之前创建的用户权限。但存在一些问题,假设用户之前的权限为: - - GRANT ALL on db1 to user; - GRANT ALL on db2 to user; - - 如果该user的dns解析后的host有10个,那么更新后的 db 权限表中,会产生 10 * 2 个条目。但我们认为当前数量级不会影响性能。 - - - -2. 自动生成的 ROOT@'%' 和 ADMIN@'%',如何更改其中的 host - - -## 权限逻辑 - -### CREATE USER - -1. create user cmy@'ip' identified by '12345'; - - 检查 cmy@'ip' 是否存在,如果存在则报错。 - -2. create user cmy@'ip' identified by '12345' default role role1; - - 检查 cmy@'ip' 是否存在,如果存在则报错。 - 赋予 cmy@'ip' role1 的所有权限,即时生效。 - role1 中加入 cmy@'ip'。 - -3. create user cmy@['domian'] identified by '12345'; - - 在白名单中检查 cmy@['domian'] 是否存在,如果存在则报错。 - -4. create user cmy@['domian'] identified by '12345' default role role1; - - 在白名单中检查 cmy@['domian'] 是否存在,如果存在则报错。 - 赋予该白名单,role1的所有权限。后台线程轮询生效。 - role1 中加入 cmy@['domian'] - -### GRANT - -1. grant select on \*.* to cmy@'ip'; - - 检查 cmy@'ip' 是否存在,不存在则报错。 - 将 select 权限赋给 cmy@'ip' on \*.* - -2. grant select on \*.* to cmy@['domain']; - - 检查 cmy@['domain'] 是否存在,不存在则报错。 - 将 select 权限赋给 cmy@['domain'] on \*.* - -3. grant select on \*.* to ROLE role1; - - 检查 role1 是否存在,不存在则报错。 - 将 select 权限赋给 role1 on \*.*。 - 将 select 权限赋给所有 role1 的用户。 - - - - - - - diff --git a/docs/design/data_export.md b/docs/design/data_export.md deleted file mode 100644 index 4e083ac2..00000000 --- a/docs/design/data_export.md +++ /dev/null @@ -1,469 +0,0 @@ -[TOC] - -# 数据导出设计 - -## 背景 - -数据导出虽然一个低频场景,但是有时却非常重要。目前palo不提供导出功能。 - -### 常见的集中导出形式 - -+ Physical(Raw) VS Logical - Physical是指以原始文件和目录的形式导出。Logical是指导出更有意义的方式,常常是文本,insert语句等。 - -+ Online VS Offline - Online是指在导出过程中不停止服务,不影响其他用户使用。注意有的时候,尤其是物理导出时,在导出过程中,需要对正在访问的物理文件加锁。 - Offline是指停止服务,进行导出。 - -+ Local VS Remote - 指相对于服务所在机器,导出到本地,还是远程机器。 - -+ Full VS Incremental - 全量导出和增量导出 - -## 需求分析 - -主要从用户的角度,主要阐述用户可能需要哪些功能,以及我们提供哪些功能。 - -### 1. 导出的数据内容? - -如下图: - - - -+ **1. 导出全量数据** -当用户希望导出的数据可以灌入其他数据库。 -另外,如果有多个表的时候,用户可能希望同时导出整个DB,或者同时导出多个表。 - -可能希望原子的导出多个表的数据,即导出的数据在同一个版本,在导出过程中,新load的数据,不被导出。 -希望原子导出:对于已经导出的数据,再次执行导出时不需要重新导出,只需要将增量部分补充即可。 - -+ **2. 按照partition导出某个表的数据** -Palo支持两层分区,尤其是对于时序性数据,历史数据分布在较老的partition,这些数据很少被使用,用户可以选择将它们导出去,然后删除这些partition,从而节省在palo中的资源占用。 - -+ **3. 任意指定数据** -用户处于临时测试等需求,可能要根据过滤条件导出一部分数据。 - -### 2. 将数据导出到哪里? - -+ 导出到用户本地 - 这种需求可以通过mysql -e "select clause" 方式进行解决。 - 这种方式最大的缺点是:如果数据量较大,常常会因为网络不稳定等原因而失败,需要多次重新执行。 -+ 导出到云存储 - - 内网用户,导出到百度内部的hadoop集群。 - - 开放云用户,导出到BOS集群。 - - 私有云用户,导出到BOS集群和社区版hadoop集群。 - -### 3. 导出数据后的访问方式? - -+ 导出的数据可以被palo快速的再次访问 - 用户可能希望导出的数据,palo还可以继续使用,只是不再关注访问的高性能,可以看做是“分级存储”功能的扩展。这种需求,用户导出后的数据作为palo的外部表存在。 -+ 导出的数据不再需要被palo访问。 - 导出数据的目的可能是希望被其他系统使用,而不再需要palo使用。 - -### 4. 导出的数据保存为什么格式? - -数据保存的格式有很多种,用户可能根据用途、使用方式、费用成本等方面的考虑,希望数据采用某种特定的格式。 -例如: -+ 为了能够快速被查看,采用csv格式。 -+ 为了能够上传和下载方便,也同时节省费用,希望导出的数据被压缩,从而减小数据大小。 -+ 为了对接HIVE等其他Hadoop生态的存储系统,希望数据直接按照HIVE的格式进行存储。 - -除此之外,用户对导出文件的数量,大小等均不是很敏感。 - -### 5. 其它 - -对于数据导出的效率,用户并不是很敏感。 - -### 需求归纳 - -用户的需求多种多样,基于Palo的定位、人力和排期,Palo不准备做一个“又大又全”的导出,而是以“解决核心需求,兼顾小众需求”的原则,对功能进行取舍。 - -这里的导出是指:“**逻辑**”导出(Logical),“**在线**导出”(Online),**全量**导出,**远程**导出。 - -详细如下: -+ 对于导出内容数据: - - 本次不支持多表同时导出,仅支持单表导出,但支持按照指定partition导出。 - - 对于任意数量的需求,我们采用"insert into 外部表 select clause"的方式解决,外部表暂时不支持分区。 - - 不支持指定版本导出,不支持增量导出。 -+ 对于导出到云存储: - 支持导出到BOS集群(云用户)和hadoop集群(内部用户) -+ 对于是否被palo访问: - 如果有被palo再次访问的需求,用户可以通过broker的方式,进行访问。 -+ 对于导出的数据格式: - 支持csv格式,并且可以指定压缩(压缩方式同broker)。不支持HIVE等其他系统的存储方式。 - -## 其他系统调研 - -### MYSQL - -#### 方法1:mysqldump命令 -``` -shell> mysqldump [arguments] > file_name - -``` - -可以通过指定导出所有数据库、部分数据库、部分表进行导出。 - -+ 如果有--tab参数,生成两个文件,一个是包含create_table语句的tab_name.sql文件,一个是纯文本文件,包含所有数据,可以指定行尾符、转义字符等。 -+ 如果没有--tab参数,生成一个文本文件。默认是生成一个SQL脚本,包含建表语句和insert语句。 - -#### 方法2:SELECT * INTO OUTFILE 'file_name' FROM tbl_name. -``` -example: - SELECT a,b,a+b - INTO OUTFILE '/tmp/result.text' - FIELDS TERMINATED BY ',' - OPTIONALLY ENCLOSED BY '"' - LINES TERMINATED BY '\n' - FROM test_table; -``` - -说明: -+ 不支持where子句,只能导出全表。 -+ **文件生成在服务端机器上,不是本地。** - -#### 方法3:INSERT ... SELECT 和 CREATE ... SELECT - -``` -Syntax: - CREATE [TEMPORARY] TABLE [IF NOT EXISTS] tbl_name [(create_definition,...)] - [table_options] - [partition_options] - [IGNORE | REPLACE] [AS] query_expression - -Syntax: - INSERT [LOW_PRIORITY | HIGH_PRIORITY] [IGNORE] [INTO] tbl_name - [PARTITION (partition_name,...)] [(col_name,...)] - SELECT ... - [ ON DUPLICATE KEY UPDATE col_name=expr, ... ] -``` - -mysql不支持外部表,所以这两种都无法将数据导出到外部。 - -### HIVE - -因为Hive的数据都是在HDFS上,所以其数据都是可以通过HDFS直接访问的。 - -#### CTAS(Create Table As Select) -``` -CREATE [EXTERNAL] TABLE table_name - [(col_name data_type [COMMENT col_comment], ... [constraint_specification])] - [PARTITIONED BY (col_name data_type, ...)] - ON ((col_value, col_value, ...), (col_value, col_value, ...), ...) - [STORED AS DIRECTORIES] - [ - [ROW FORMAT row_format] - [STORED AS file_format] - | STORED BY 'storage.handler.class.name' [WITH SERDEPROPERTIES (...)] -- (Note: Available in Hive 0.6.0 and later) - ] - [LOCATION hdfs_path] - - [AS select_statement]; -- (Note: Available in Hive 0.5.0 and later; not supported for external tables) -``` - -+ 只有当select执行完毕,并导入以后,该表才可见。 -+ 如果希望将select语句的结果导入到hdfs_path中,将location指定为hdfs_path即可。 -+ 有一些限制: - - The target table cannot be a partitioned table. - - The target table cannot be an external table. - - The target table cannot be a list bucketing table. -+ 如果指定了数据的存储格式,那么会自动被转化。 - -例如: -``` -CREATE TABLE tbl_text - ROW FORMAT DELIMITED FIELDS TERMINATED BY ',' - LOCATION '/tmp/data' - AS select * from tbl where condition; -``` - -#### Writing data into the filesystem from queries - -```sql -Standard syntax: -INSERT OVERWRITE [LOCAL] DIRECTORY directory1 - [ROW FORMAT row_format] [STORED AS file_format] - SELECT ... FROM ... - -Hive extension (multiple inserts): -FROM from_statement -INSERT OVERWRITE [LOCAL] DIRECTORY directory1 select_statement1 -[INSERT OVERWRITE [LOCAL] DIRECTORY directory2 select_statement2] ... - - -row_format - : DELIMITED [FIELDS TERMINATED BY char [ESCAPED BY char]] [COLLECTION ITEMS TERMINATED BY char] - [MAP KEYS TERMINATED BY char] [LINES TERMINATED BY char] - [NULL DEFINED AS char] (Note: Only available starting with Hive 0.13) -``` -说明: -+ 如果指定了LOCAL关键字,那么数据写入到本地文件系统。否则,写入到HDFS中。 -+ 文件格式是文本文件,多个列之间用^A字符('\001')分隔。对于非Primitive type的列,转化为json格式。 -+ 从0.11版本之后,列分隔符可以指定。 -+ OVERWRITE关键字是必须的。如果指定的目录已经存在,那么会被覆盖。 -+ 在一条语句中,可以指定多个INSERT子句。可以同时包含HDFS目录,local目录,以及table(或者partition)。 -+ 对于大量数据,使用“INSERT OVERWRITE statements to HDFS filesystem directories”是最好的方式。因为这时,Hive可以在Map-Reduce作业中并行写入。 -+ 在一个语句中指定多个insert,可以减少原始数据被扫描的次数。 - - -#### Inserting values into tables from SQL -``` -Standard syntax: - INSERT OVERWRITE TABLE tablename1 [PARTITION (partcol1=val1, partcol2=val2 ...) [IF NOT EXISTS]] select_statement1 FROM from_statement; - INSERT INTO TABLE tablename1 [PARTITION (partcol1=val1, partcol2=val2 ...)] select_statement1 FROM from_statement; - -Hive extension (multiple inserts): - FROM from_statement - INSERT OVERWRITE TABLE tablename1 [PARTITION (partcol1=val1, partcol2=val2 ...) [IF NOT EXISTS]] select_statement1 - [INSERT OVERWRITE TABLE tablename2 [PARTITION ... [IF NOT EXISTS]] select_statement2] - [INSERT INTO TABLE tablename2 [PARTITION ...] select_statement2] ...; - - FROM from_statement - INSERT INTO TABLE tablename1 [PARTITION (partcol1=val1, partcol2=val2 ...)] select_statement1 - [INSERT INTO TABLE tablename2 [PARTITION ...] select_statement2] - [INSERT OVERWRITE TABLE tablename2 [PARTITION ... [IF NOT EXISTS]] select_statement2] ...; - -Hive extension (dynamic partition inserts): -INSERT OVERWRITE TABLE tablename PARTITION (partcol1[=val1], partcol2[=val2] ...) select_statement FROM from_statement; -INSERT INTO TABLE tablename PARTITION (partcol1[=val1], partcol2[=val2] ...) select_statement FROM from_statement; -``` - -说明: -+ 上述三种方式中,导出的文件数是和mapper和reducer的个数相关的,如果没有reduce,则文件数为mapper的个数;否则文件数为reducer的个数。 - -#### EXPORT TABLE - -``` -Export Syntax: - EXPORT TABLE tablename [PARTITION (part_column="value"[, ...])] - TO 'export_target_path' [ FOR replication('eventid') ] -``` -+ Export不仅是data还有metadata。 metadata放在目标dir中,数据放在子目录中。 -+ 数据的组织形式和原table的形式一致,原table有多少个文件,export出的就有多少文件。 -+ 可以导出一整张表,也可以导出指定的partition。 - -### GreenplumDB - -#### 方法1:INSERT ... SELECT - -``` -example: - CREATE WRITABLE EXTERNAL TABLE table_name - LOCATION ('gphdfs://hdfshost:port/path') - FORMAT format - DISTRIBUTED BY distributed_key; - - INSERT INTO table_name_1 SELECT * FROM table_name_2; -``` - -+ 通过外部表导出可以并行导出,并发数为Segment Host的数目。 - -#### 方法2:COPY - -``` -Syntax: - COPY (select_statement) TO 'target_path'; -``` - -+ 数据需要在master上汇总,因此不能并行。 -+ 这种方式的导出没有方法1高,但是其简单方便,主要用于数据量较小的情况。 - -### VerticaDB - -#### 方法1:通过S3EXPORT()导出到S3 - -``` -example: - SELECT S3EXPORT(* USING PARAMETERS url='s3://exampleBucket/object') FROM exampleTable; -``` - -#### 方法2:EXPORT导出到另一个Vertica DB中 - -``` -example1:导出整张表 - EXPORT TO VERTICA testdb.customer_dimension FROM customer_dimension; - -example2:导出过滤后的表 - EXPORT TO VERTICA testdb.ma_customers - AS SELECT customer_key, customer_name, annual_income - FROM customer_dimension - WHERE customer_state = 'MA'; - -example3: 导出部分列 - EXPORT TO VERTICA testdb.people(name, gender, age) - FROM customer_dimension (customer_name, customer_gender, customer_age); -``` - -### Impala - -Impala虽然支持Hive的元数据,语法格式也与Hive内饰。但是并不支持Hive的EXPORT TABLE语法。目前只支持以external table的方式导出数据,包括CTAS、insert ... select ... 的方式。 - -``` -Syntax: - INSERT { INTO | OVERWRITE } [TABLE] table_name - [(column_list)] - [ PARTITION (partition_clause)] - select_statement - - - partition_clause ::= col_name [= constant] [, col_name [= constant] ...] - - hint_clause ::= [SHUFFLE] | [NOSHUFFLE] (Note: the square brackets are - part of the syntax.) -``` - -说明: -+ 尽管可以在select语句中指定order by,但是会被忽略。因为数据会在多台机器上进行,每个机器会生成一个文件,所以将数据排序没有实际意义。 -+ 任何一个insert语句都会生成不同的名字的文件,所以多个insert into ...select语句可以同时执行(注意不是 INSERT OVERWRITE语句)。 -+ 默认会在目标目录下生成一个以"_dir"记为的子目录,作为临时目录。在执行完毕后,执行mv操作。 - -## 概要设计 - -### API - -采用类似于Hive的方式: -两种方式: - -``` -Syntax: - -EXPORT TABLE tablename - [PARTITION (name1[, ...])] - TO 'export_target_path' - [PROPERTIES("key"="value")] - BY BROKER 'broker_name' [( $broker_attrs)] - -INSERT { INTO | OVERWRITE } TABLE table_name[(column_list)] - select_statement - [BY BROKER 'broker_name' [( $broker_attrs)]] - -``` - -对于export语句,使用的BROKER通过by子句进行指定。行分隔符、压缩方式使用PROPERTIES子句指定。 - -insert...select语句中,并且需要使用建表时指定的broker,那么可以不提供BY BROKER子句,否则需要重新指定broker。行分隔符等属性,采用建外部表时指定的值,如未指定则取默认值。 - -### 导出的数据格式 - -导出数据为CSV格式,可以指定列分隔符; - -支持对数据压缩,支持的压缩方式同broker scan。 - -### 异步执行 - -数据导出一般执行时间较长,所以采用异步方式,用户提交命令后,会先返回,然后后台执行,同时提供方式供用户查看执行进度。 - -和现有insert ... select 语句类似,对外部表的insert ..。 select 语句仍然使用**SHOW LOAD**来查看进度。 - -对于EXPORT语句,使用**SHOW EXPORT**查看进度。 - -考虑到排期,本期优先实现insert external table select...。 - -### 持久化 - -因为是异步执行,为了防止Palo集群宕机后,用户看不到export任务,所以需要对导出任务进行持久化。 - -### 导出任务的状态切换 - -和现有load作业的状态切换类似。 - - - -### SHOW LOAD - -兼容现有show load语句。 - -### SHOW EXPORT -和SHOW LOAD类似。 - -``` - SHOW EXPORT - [FROM db_name] - [ - WHERE - [LABEL [ = "your_label" | LIKE "label_matcher"]] - [STATE = ["PENDING"||"EXPORTING"|"FINISHED"|"CANCELLED"|]] - ] - [ORDER BY ...] - [LIMIT limit]; -``` - -### 执行计划 - -#### export语句的执行计划 - -如下图: - - - -#### insert ... select语句的执行计划 - - - -### 并行导出 - -Palo中的外部表,**不支持多分区方式。** - -对于源数据表多partition的表,可以做到多个BE同时并行的export,从而提高效率。 - -### 需要snapshot - -对于export任务,因为持续时间较长,为了避免在导出过程中be将数据merge掉,导致找不到对应的数据。所以需要在执行之前要对BE上的tablet打一个snapshot,同时在任务结束后,及时将snapshot删除。 - -### 资源隔离 - -这里主要是指IO操作的隔离。 - -Export操作是一个**低频、低优先级**的操作,如果按照现有的查询逻辑进行自行,在Palo当前的资源隔离实现时,那么当表中数据量很大时,可能会将BE扫描线程打满,影响到正常的查询服务,所以需要对export的操作进行限速。 - -鉴于Palo后面会在底层进行IO隔离,不同用户间按比例调度,同一个用户按round robin进行调度,并且内部实现类似时间片的分配机制。这样即便export时查询很大,可以保证不会影响到其他用户,以及同一用户的其他查询。 - -在目前,采用的隔离方案: -+ 对于Export方式,对于多个Fragment instance,我们不在使用MPP的执行方式,而是采用串行的执行方式。(因为只有一级scan node,并且不需要节点进行汇聚数据,所以可以使用串行) -+ 对于insert select方式,因为select语句可能比较复杂,例如包含join等其他操作,不能采用串行的方式。 - -作为临时方案,为了简单起见,对于同一个fragment,执行时仍然是“每个BE只会有一个Fragment instance”。这样,虽然不会因为export任务将所有BE都打满,但是在极端情况下,仍然可能会将单台BE打满。但考虑到该功能主要在公有云上使用,每个用户都是自己独立的集群,并且是一个低频的操作,同时待新的IO隔离方式上线后,该功能就不需要了。 - -### 文件目录 - -目前Broker scan时,暂不支持子目录,所以导出的数据都在同一个目录下,不生成子目录。 - -多个BE并行导出,即多个BE同时向一个目录下写数据(写入到多个文件)。 - -**写入顺序**:先在临时目录下写入,待写入完成后,mv到目标目录。 - -这里要注意的是:在极端情况下,导出任务失败的情况下,例如Palo集群宕机等,可能会将数据遗留在临时目录中。这时需要手动进行清理。(hdfs dfs -rm -r tmp_dir) - -### 文件命名格式: - -每个tablet默认对应一个文件,如果文件较大时,会进行拆分,通过后缀的递增进行区分。 - -``` -文件命名格式: - [$partition_id_]{$instance_id}_{$random_int}.data.{$increasing_int} -``` - -### 文件大小 - -在broker scan时,对于单个大文件,目前只支持文本文件的拆分读取。如果是压缩文件,无法做到拆分读取。 - -同时在向BOS写入时,目前是要先写本地文件,然后在同一向BOS上传,如果文件过大,失败的可能性比较大。同时目前broker本身暂不支持自动将文件拆分写入。 - -综上,需要在上层对文件大小进行限制,尤其对压缩的HDFS文件和输出到BOS时。 - -初期可以采用比较简单的分拆方法,比如行数或者压缩前buffer大小。 - -### 错误重试 - -对于export方式,每台BE上的fragment可以单独进行重试。 - -对于insert...select方式,和普通查询一样,不进行重试。 - -### 忽略数据排序 - -对于insert select语句,select子句中可能会有order by,但是会被忽略。因为数据会在多台机器上进行,每个机器会生成一个文件,所以将数据排序没有实际意义。 - -### 不支持增量导出 - -即连续的两次导出任务,不支持用户在第二次仅导出第一次没有导出的数据。 diff --git a/docs/design/metadata_design.md b/docs/design/metadata_design.md deleted file mode 100644 index 61d22cf3..00000000 --- a/docs/design/metadata_design.md +++ /dev/null @@ -1,102 +0,0 @@ -# Palo 元数据设计文档 - -## 名词解释 - -* FE:Frontend,即 Palo 的前端节点。主要负责接收和返回客户端请求、元数据以及集群管理、查询计划生成等工作。 -* BE:Backend,即 Palo 的后端节点。主要负责数据存储与管理、查询计划执行等工作。 - -## 整体架构 - - - -如上图,Palo 的整体架构分为两层。多个 FE 组成第一层,提供 FE 的横向扩展和高可用。多个 BE 组成第二层,负责数据存储于管理。本文主要介绍 FE 这一层中,元数据的设计与实现方式。 - -1. FE 节点分为 follower 和 observer 两类。各个 FE 之间,通过 bdbje([BerkeleyDB Java Edition](http://www.oracle.com/technetwork/database/database-technologies/berkeleydb/overview/index-093405.html))进行 leader 选举,数据同步等工作。 - -2. follower 节点通过选举,其中一个 follower 成为 leader 节点,负责元数据的写入操作。当 leader 节点宕机后,其他 follower 节点会重新选举出一个 leader,保证服务的高可用。 - -3. observer 节点仅从 leader 节点进行元数据同步,不参与选举。可以横向扩展以提供元数据的读服务的扩展性。 - -> 注:follower 和 observer 对应 bdbje 中的概念为 replica 和 observer。下文可能会同时使用两种名称。 - -## 元数据结构 - -Palo 的元数据是全内存的。每个 FE 内存中,都维护一个完整的元数据镜像。在百度内部,一个包含2500张表,100万个分片(300万副本)的集群,元数据在内存中仅占用约 2GB。(当然,查询所使用的中间对象、各种作业信息等内存开销,需要根据实际情况估算。但总体依然维持在一个较低的内存开销范围内。) - -同时,元数据在内存中整体采用树状的层级结构存储,并且通过添加辅助结构,能够快速访问各个层级的元数据信息。 - -下图是 Palo 元信息所存储的内容。 - - - -如上图,Palo 的元数据主要存储4类数据: - -1. 用户数据信息。包括数据库、表的Schema、分片信息等。 -2. 各类作业信息。如导入作业,Clone作业、SchemaChange作业等。 -3. 用户及权限信息。 -4. 集群及节点信息。 - -## 数据流 - - -元数据的数据流具体过程如下: - -1. 只有 leader FE 可以对元数据进行写操作。写操作在修改 leader 的内存后,会序列化为一条log,按照 key-value 的形式写入 bdbje。其中 key 为连续的整型,作为 log id,value 即为序列化后的操作日志。 - -2. 日志写入 bdbje 后,bdbje 会根据策略(写多数/全写),将日志复制到其他 non-leader 的 FE 节点。non-leader FE 节点通过对日志回放,修改自身的元数据内存镜像,完成与 leader 节点的元数据同步。 - -3. leader 节点的日志条数达到阈值后(默认 10w 条),会启动 checkpoint 线程。checkpoint 会读取已有的 image 文件,和其之后的日志,重新在内存中回放出一份新的元数据镜像副本。然后将该副本写入到磁盘,形成一个新的 image。之所以是重新生成一份镜像副本,而不是将已有镜像写成 image,主要是考虑写 image 加读锁期间,会阻塞写操作。所以每次 checkpoint 会占用双倍内存空间。 - -4. image 文件生成后,leader 节点会通知其他 non-leader 节点新的 image 已生成。non-leader 主动通过 http 拉取最新的 image 文件,来更换本地的旧文件。 - -5. bddje 中的日志,在 image 做完后,会定期删除旧的日志。 - -## 实现细节 - -### 元数据目录 - -1. 元数据目录通过 FE 的配置项 `meta_dir` 指定。 - -2. `bdb/` 目录下为 bdbje 的数据存放目录。 - -3. `image/` 目录下为 image 文件的存放目录。 - - * `image.[logid]` 是最新的 image 文件。后缀 `logid` 表明 image 所包含的最后一条日志的 id。 - * `image.ckpt` 是正在写入的 image 文件,如果写入成功,会重命名为 `image.[logid]`,并替换掉就的 image 文件。 - * `VERSION` 文件中记录着 `cluster_id`。`cluster_id` 唯一标识一个 Palo 集群。是在 leader 第一次启动时随机生成的一个 32 位整型。也可以通过 fe 配置项 `cluster_id` 来指定一个 cluster id。 - * `ROLE` 文件中记录的 FE 自身的角色。只有 `FOLLOWER` 和 `OBSERVER` 两种。其中 `FOLLOWER` 表示 FE 为一个可选举的节点。(注意:即使是 leader 节点,其角色也为 `FOLLOWER`) - -### 启动流程 - -1. FE 第一次启动,如果启动脚本不加任何参数,则会尝试以 leader 的身份启动。在 FE 启动日志中会最终看到 `transfer from UNKNOWN to MASTER`。 - -2. FE 第一次启动,如果启动脚本中指定了 `-helper` 参数,并且指向了正确的 leader FE 节点,那么该 FE 首先会通过 http 向 leader 节点询问自身的角色(即 ROLE)和 cluster_id。然后拉取最新的 image 文件。读取 image 文件,生成元数据镜像后,启动 bdbje,开始进行 bdbje 日志同步。同步完成后,开始回放 bdbje 中,image 文件之后的日志,完成最终的元数据镜像生成。 - - > 注1:使用 `-helper` 参数启动时,需要首先通过 mysql 命令,通过 leader 来添加该 FE,否则,启动时会报错。 - - > 注2:`-helper` 可以指向任何一个 follower 节点,即使它不是 leader。 - - > 注2:bdbje 在同步日志过程中,fe 日志会显示 `xxx detached`, 此时正在进行日志拉取,属于正常现象。 - -3. FE 非第一次启动,如果启动脚本不加任何参数,则会根据本地存储的 ROLE 信息,来确定自己的身份。同时根据本地 bdbje 中存储的集群信息,获取 leader 的信息。然后读取本地的 image 文件,以及 bdbje 中的日志,完成元数据镜像生成。(如果本地 ROLE 中记录的角色和 bdbje 中记录的不一致,则会报错。) - -4. FE 非第一次启动,且启动脚本中指定了 `-helper` 参数。则和第一次启动的流程一样,也会先去询问 leader 角色。但是会和自身存储的 ROLE 进行比较。如果不一致,则会报错。 - -#### 元数据读写与同步 - -1. 用户可以使用 mysql 连接任意一个 FE 节点进行元数据的读写访问。如果连接的是 non-leader 节点,则该节点会将写操作转发给 leader 节点。leader 写成功后,会返回一个 leader 当前最新的 log id。之后,non-leader 节点会等待自身回放的 log id 大于回传的 log id 后,才将命令成功的消息返回给客户端。这种方式保证了任意 FE 节点的 Read-Your-Write 语义。 - - > 注:一些非写操作,也会转发给 leader 执行。比如 `SHOW LOAD` 操作。因为这些命令通常需要读取一些作业的中间状态,而这些中间状态是不写 bdbje 的,因此 non-leader 节点的内存中,是没有这些中间状态的。(FE 直接的元数据同步完全依赖 bdbje 的日志回放,如果一个元数据修改操作不写 bdbje 日志,则在其他 non-leader 节点中是看不到该操作修改后的结果的。) - -2. leader 节点会启动一个 TimePrinter 线程。该线程会定期向 bdbje 中写入一个当前时间的 key-value 条目。其余 non-leader 节点通过回放这条日志,读取日志中记录的时间,和本地时间进行比较,如果发现和本地时间的落后大于指定的阈值(配置项:`meta_delay_toleration_second`。写入间隔为该配置项的一半),则该节点会处于**不可读**的状态。此机制解决了 non-leader 节点在长时间和 leader 失联后,仍然提供过期的元数据服务的问题。 - -3. 各个 FE 的元数据只保证最终一致性。正常情况下,不一致的窗口期仅为毫秒级。我们保证同一 session 中,元数据访问的单调一致性。但是如果同一 client 连接不同 FE,则可能出现元数据回退的现象。(但对于批量更新系统,该问题影响很小。) - -### 宕机恢复 - -1. leader 节点宕机后,其余 follower 会立即选举出一个新的 leader 节点提供服务。 -2. 当多数 follower 节点宕机时,元数据不可写入。当元数据处于不可写入状态下,如果这时发生写操作请求,目前的处理流程是 **FE 进程直接退出**。后续会优化这个逻辑,在不可写状态下,依然提供读服务。 -3. observer 节点宕机,不会影响任何其他节点的状态。也不会影响元数据在其他节点的读写。 - - - diff --git a/docs/design/multi_tenant.md b/docs/design/multi_tenant.md deleted file mode 100644 index b3f04f7f..00000000 --- a/docs/design/multi_tenant.md +++ /dev/null @@ -1,212 +0,0 @@ -# 多租户设计 - -## 背景 -palo作为一款PB级别的在线报表与多维分析数据库,对外通过开放云提供云端的数据库服务,并且对于多租户部署了多套物理集群。对内,一套物理集群部署了多个业务,对于隔离性要求比较高的业务单独搭建了集群。针对以上存在几点问题: - -- 部署多套物理集群维护代价大(升级、功能上线、bug修复)。 -- 一个用户的查询或者查询引起的bug经常会影响其他用户。 -- 实际生产环境单机只能部署一个BE进程。而多个BE可以更好的解决胖节点问题。并且对于join、聚合操作可以提供更高的并发度。 - -综合以上三点,palo需要新的多租户方案,既能做到较好的资源隔离和故障隔离,同时也能减少维护的代价,满足共有云和私有云的需求。 - -## 设计原则 - -- 使用简单 -- 开发代价小 -- 方便现有集群的迁移 - -## 名词解释 - -- FE: Frontend,即 Palo 中用于元数据管理即查询规划的模块。 -- BE: Backend,即 Palo 中用于存储和查询数据的模块。 -- Master: FE 的一种角色。一个Palo集群只有一个Master,其他的FE为Observer或者Follower。 -- instance:一个 BE 进程及时一个 instance。 -- host:单个物理机 -- cluster:即一个集群,由多个instance组成。 -- 租户:一个cluster属于一个租户。cluster和租户之间是一对一关系。 -- database:一个用户创建的数据库 - -## 主要思路 - -- 一个host上部署多个BE的instance,在进程级别做资源隔离。 -- 多个instance形成一个cluster,一个cluster分配给一个业务独立的的租户。 -- FE增加cluster这一级并负责cluster的管理。 -- CPU,IO,内存等资源隔离采用cgroup。 - -## 设计方案 - -为了能够达到隔离的目的,引入了**虚拟cluster**的概念。 - -1. cluster表示一个虚拟的集群,由多个BE的instance组成。多个cluster共享FE。 -2. 一个host上可以启动多个instance。cluster创建时,选取任意指定数量的instance,组成一个cluster。 -3. 创建cluster的同时,会创建一个名为superuser的账户,隶属于该cluster。superuser可以对cluster进行管理、创建数据库、分配权限等。 -4. Palo启动后,汇创建一个默认的cluster:default_cluster。如果用户不希望使用多cluster的功能,则会提供这个默认的cluster,并隐藏多cluster的其他操作细节。 - -具体架构如下图: - - - - -## SQL 接口 - -- 登录 - - 默认集群登录名: user_name@default_cluster 或者 user_name - - 自定义集群登录名:user_name@cluster_name - - `mysqlclient -h host -P port -u user_name@cluster_name -p password` - -- 添加、删除、下线(decommission)以及取消下线BE - - `ALTER SYSTEM ADD BACKEND "host:port"` - `ALTER SYSTEM DROP BACKEND "host:port"` - `ALTER SYSTEM DECOMMISSION BACKEND "host:port"` - `CANCEL DECOMMISSION BACKEND "host:port"` - - 强烈建议使用 DECOMMISSION 而不是 DROP 来删除 BACKEND。DECOMMISSION 操作会首先将需要下线节点上的数据拷贝到集群内其他instance上。之后,才会真正下线。 - -- 创建集群,并指定superuser账户的密码 - - `CREATE CLUSTER cluster_name PROPERTIES ("instance_num" = "10") identified by "password"` - -- 进入一个集群 - - `ENTER cluster_name` - -- 集群扩容、缩容 - - `ALTER CLUSTER MODIFY cluster_name PROPERTIES ("instance_num" = "10")` - - 当指定的实例个数多于cluster现有be的个数,则为扩容,如果少于则为缩容。 - -- 链接、迁移db - - `LINK DATABASE src_cluster_name.db_name dest_cluster_name.db_name` - - 软链一个cluster的db到另外一个cluster的db ,对于需要临时访问其他cluster的db却不需要进行实际数据迁移的用户可以采用这种方式。 - - `MIGRATE DATABASE src_cluster_name.db_name dest_cluster_name.db_name` - - 如果需要对db进行跨cluster的迁移,在链接之后,执行migrate对数据进行实际的迁移。 - - 迁移不影响当前两个db的查询、导入等操作,这是一个异步的操作,可以通过`SHOW MIGRATIONS`查看迁移的进度。 - -- 删除集群 - - `DROP CLUSTER cluster_name` - - 删除集群,要求先手动删除的集群内所有database。 - -- 其他 - - `SHOW CLUSTERS` - - 展示系统内已经创建的集群。只有root用户有该权限。 - - `SHOW BACKENDS` - - 查看集群内的BE instance。 - - `SHOW MIGRATIONS` - - 展示当前正在进行的db迁移任务。执行完db的迁移后可以通过此命令查看迁移的进度。 - -## 详细设计 - -1. 命名空间隔离 - - 为了引入多租户,需要对系统内的cluster之间的命名空间进行隔离。 - - palo现有的元数据采用的是image + journal 的方式(元数据的设计见相关文档)。palo会把涉及元数据的操作的记录为一个 journal (操作日志),然后定时的按照**图1**的方式写成image,加载的时候按照写入的顺序读即可。但是这样就带来一个问题已经写入的格式不容易修改,比如记录数据分布的元数据格式为:database+table+tablet+replica 嵌套,如果按照以往的方式要做cluster之间的命名空间隔离,则需要在database上增加一层cluster,内部元数据的层级变为:cluster+database+table+tablet+replica,如**图2**所示。但加一层带来的问题有: - - - 增加一层带来的元数据改动,不兼容,需要按照图2的方式cluster+db+table+tablet+replica层级写,这样就改变了以往的元数据组织方式,老版本的升级会比较麻烦,比较理想的方式是按照图3在现有元数据的格式下顺序写入cluster的元数据。 - - - 代码里所有用到db、user等,都需要加一层cluster,一工作量大改动的地方多,层级深,多数代码都获取db,现有功能几乎都要改一遍,并且需要在db的锁的基础上嵌套一层cluster的锁。 - -  - - 综上这里采用了一种通过给db、user名加前缀的方式去隔离内部因为cluster之间db、user名字冲突的问题。 - - 如下,所有的sql输入涉及db名、user名的,都需要根据自己所在的cluster来拼写db、user的全名。 - -  - - 采用这种方式以上两个问题不再有。元数据的组织方式也比较简单。即采用**图3**每个cluster记录下属于自己cluster的db、user,以及节点即可。 - -2. BE 节点管理 - - 每个cluster都有属于自己的一组instance,可以通过`SHOW BACKENDS`查看,为了区分出instance属于哪个cluster以及使用情况,BE引入了多个状态: - - - free:当一个BE节点被加入系统内,此时be不属于任何cluster的时候处于空闲状态 - - using:当创建集群、或者扩容被选取到一个cluster内则处于使用中。 - - cluster decommission:如果执行缩容量,则正在执行缩容的be处于此状态。结束后,be状态变为free。 - - system decommission:be正在下线中。下线完成后,该be将会被永久删除。 - - 只有root用户可以通过`SHOW PROC "/backends"`中cluster这一项查看集群内所有be的是否被使用。为空则为空闲,否则为使用中。`SHOW BACKENDS`只能看到所在cluster的节点。以下是be节点状态变化的示意图。 - -  - -3. 创建集群 - - 只有root用户可以创建一个cluster,并指定任意数量的BE instance。 - - 支持在相同机器上选取多个instance。选择instance的大致原则是:尽可能选取不同机器上的be并且使所有机器上使用的be数尽可能均匀。 - - 对于使用来讲,每一个user、db都属于一个cluster(root除外)。为了创建user、db,首先需要进入一个cluster。在创建cluster的时候系统会默认生成这个cluster的管理员,即superuser账户。superuser具有在所属cluster内创建db、user,以及查看be节点数的权限。所有的非root用户登录必须指定一个cluster,即`user_name@cluster_name`。 - - 只有root用户可以通过`SHOW CLUSTER`查看系统内所有的cluster,并且可以通过@不同的集群名来进入不同的cluster。对于除了root之外的用户cluster都是不可见的。 - - 为了兼容老版本palo内置了一个名字叫做default_cluster的集群,这个名字在创建集群的时候不能使用。 - -  - -4. 集群扩容 - - 集群扩容的流程同创建集群。会优先选取不在集群之外的host上的BE instance。选取的原则同创建集群。 - -5. 集群缩容、CLUSTER DECOMMISSION - - 用户可以通过设置 cluster 的 instance num 来进行集群缩容。 - - 集群的缩容会优先在BE instance 数量最多的 host 上选取 instance 进行下线。 - - 用户也可以直接使用 `ALTER CLUSTER DECOMMISSION BACKEND` 来指定BE,进行集群缩容。 - -  - -6. 建表 - - 为了保证高可用,每个分片的副本必需在不同的机器上。所以建表时,选择副本所在be的策略为在每个host上随机选取一个be。然后从这些be中随机选取所需副本数量的be。总体上做到每个机器上分片分布均匀。 - - 因此,加入需要创建一个3副本的分片,即使cluster包含3个或以上的instance,但是只有2个或以下的host,依然不能创建该分片。 - -7. 负载均衡 - - 负载均衡的粒度为cluster级别,cluster之间不做负载均衡。但是在计算负载是在host一级进行的,而一个host上可能存在多个不同cluster的BE instance。 cluster内,会通过每个host上所有分片数目、存储使用率计算负载,然后把负载高的机器上的分片往负载低的机器上拷贝(详见负载均衡相关文档)。 - -8. LINK DATABASE(软链) - - 多个集群之间可以通过软链的方式访问彼此的数据。链接的级别为不同cluster的db。 - - 通过在一个cluster内,添加需要访问的其他cluster的db的信息,来访问其他cluster中的db。 - - 当查询链接的db时,所使用的计算以及存储资源为源db所在cluster的资源。 - - 被软链的db不能在源cluster中删除。只有链接的db被删除后,才可以删除源db。而删除链接db,不会删除源db。 - -9. MIGRATE DATABASE - - db可以在cluster之间进行物理迁移。 - - 要迁移db,必须先链接db。执行迁移后数据会迁移到链接的db所在的cluster,并且执行迁移后源db被删除,链接断开。 - - 数据的迁移,复用了负载均衡以及副本恢复中,复制数据的流程(详见负载均衡相关文档)。具体实现上,在执行`MIRAGTE`命令后,Palo会在元数据中,将源db的所有副本所属的cluster,修改为目的cluster。 - - Palo会定期检查集群内机器之间是否均衡、副本是否齐全、是否有多余的副本。db的迁移即借用了这个流程,在检查副本齐全的时候同时检查副本所在的be是否属于该cluster,如果不属于,则记入要恢复的副本。并且副本多余要删除的时候会优先删除cluster外的副本,然后再按照现有的策略选择:宕机的be的副本->clone的副本->版本落后的副本->负载高的host上的副本,直到副本没有多余。 - -  - -10. BE的进程隔离 - - 为了实现be进程之间实际cpu、io以及内存的隔离,需要依赖于be的部署。部署的时候需要在外围配置cgroup,把要部署的be的进程都写入cgroup。如果要实现io的物理隔离各be配置的数据存放路径需要在不同磁盘上,这里不做过多的介绍。 diff --git a/docs/design/real_time_loading.md b/docs/design/real_time_loading.md deleted file mode 100644 index ef68d5a3..00000000 --- a/docs/design/real_time_loading.md +++ /dev/null @@ -1,206 +0,0 @@ -# 一些名词 - -- Transaction: 在描述中事务和导入是同一个概念,为了方便以后事务的集成,这里设计框架的时候考虑了事务的内容,而导入是事务的子内容,所以很多名词直接用事务了。 - -# 数据结构 -## FE 中的数据结构 -### 事务的状态 -TransactionState 有4个状态: - -- PREPARE 表示事务已经开始了,Coordinator已经开始导入数据了。对于我们改造的旧系统来说就是Loading阶段开始了,由FE启动一个模块当做Coordinator来发起任务。 -- COMMITTED 表示事务已经结束了,此时数据在BE上已经quorum存在了,但是FE并没有通知BE这个导入的version是多少,所以数据对查询还不可见。 -- VISIBLE 表示这次导入事务对查询可见了。 -- ABORTED 表示导入由于内部故障导致失败或者由用户显示的取消了导入任务,由reason字段来说明具体的原因 - -### 事务表相关的结构 - -table_version 表 - -table_id | visible_version | next_version ----|---|--- -1001 | 11 | 20 -1002 | 13 | 19 - -- visible_version 表示这个table可见的版本号,也就是查询时给查询计划指定的版本号 -- next_version 是version通知机制运行时给一个导入事务分配的版本号,分配完之后自动+1 - -transaction_state 表 - -transaction_id | affected_table_version | coordinator | transaction_state |start_time| commit_time | finish_time ----|---|---|---|--|--|--| - 11 | <1001, 12><1002,13> - 13 | <1002,14> - -- transactiond_id 表示一个导入事务唯一的ID -- affected_table_version 表示一个导入事务涉及到的所有table,这个列表对检测导入冲突有用。这里记录的是一个tuple list,<tableid, version> 的list,在publishVersion的时候会用。 -- coordinator 记录当前的协调节点,为以后的实时导入服务的,目前就是FE -- start_time 是导入事务进入prepare 状态的时间 -- commit_time 是load完成进入committed状态的时间 -- finish_time 是一个事务完成的时间,可以是visible的时间,可以是abort的时间 -- 这里没有记录label,需要在现有的导入元数据中记录一个transactionid,把label映射到实时导入框架的transactionid上 - -## BE 中的数据结构 -### BE 文件目录结构 - -目前没有看过BE的目录结构,这里的只是原理上的设计,我们以后可以改,但是思路是这样的。 假定每个tablet都有一个目录,目录结构如下: - -- data 存储的就是现在的tablet的数据。 -- staging 是数据刚导入的时候把数据存放在这个目录下,等修改version的时候把数据再迁移到deltas和data目录下,在version通知机制的api里会详细写一下。 -- delta 存储的是delta文件,这个文件不会被compaction,会等待一定的时间(比如30 min后自动删除),这个文件夹主要用于副本恢复使用。 - -# API 设计 -## FE 实时导入相关API - - service GlobalTransactionManager { - // 这个API 供Coordinator发起导入任务时调用, 这个API的预期行为是: - // 1 递增的生成一个TransactionID - // 2 当前发起请求的节点作为Coordinator节点 - // 3 事务的初始状态为PREPARE - // 4 把TransactionID, Label,Coordinator, TransactionState 信息写入事务表,TableList和CommitTime为空 - // 3. 把TransactionID信息返回给调用的Coordinator节点 - int64 beginTransaction(); - - // 当coordinator导入完成后,把导入的完成情况汇报给FE时使用 - // label 是导入时用户指定的唯一标识 - // tabletsStatus 表示的是本次导入过程中涉及到的各个table的各个分片的副本的导入情况,status表示导入成功还是失败 - // fe收到这个信息后,做以下两个动作: - // 1. 判断这次导入是否成功,这里主要是一致性方面,如果根据一致性协议判定导入失败,那么就在事务表中把事务状态标记为FAILED,然后返回客户端失败; 另外也需要判断这个导入是否被CANCEL,如果已经CANCEL直接返回失败信息即可。 - // 2. 如果判定导入成功,那么 - // 2.1 将导入失败的tablet标记为CLONE状态,然后生成RepairTabletTask,如果RepairTabletTask的version比现在commit的version小的话,那么要生成一个新的task,然后下发给BE - // 2.2 计算本次导入涉及到的所有的tableId,将tableid的信息写入事务表,在事务表中把事务的状态标记为COMMITTED - Status commitTransaction(int64 transactionId, list<tuple<tabletId,Status>> tabletsStatus) - - Status rollbackTransaction(int64 transactionId) - - // 获取一个table的一个transaction对应的version - map<transactionid, version> getTransactionVersion(list<tuple<TransactionIds, tableId>>) - } - - // 这个模块用来充当DPP,小批量,BROKER方式导入的Coordinator - Service TransactionCoordinator { - // BE 执行完load命令后,调用这个API给FE汇报导入完成情况,FE收到这个消息后,在内存中保存下来 - status finishLoad(list<tuple<transactionid,tabletid,status>>) - } - - Service TabletReplicationManager { - status finishClone() - } - -## BE Agent的API - struct RepairTabletTask { - // 表示从哪个BE上来读取数据 - string sourceIp; - // 表示当前故障的BE的IP - string targetIp; - // 修复数据的version号 - string version; - } - - service LocalTransactonManager { - // FE 调用BE,让BE从sourceBE上clone数据, version表示clone到哪个版本的数据,这个API注意以下几点: - // 1. BE 在执行时需要跟sourceIp比对一下需要传递哪些数据,哪些数据本地有,哪些数据本地没有,是否需要全量恢复。 - // 2. sourceIP 对应的BE节点上version的数据可能还暂时不存在,那么BE需要等待一下。 - // 3. 这个API 要实现成为一个能够增量执行的API,因为FE会不断的发送不同的Repair任务给BE,比如FE发送一个version=8的,后来又发送一个version=9的,那么BE需要判断本地是否有repair任务在执行,如果在执行,那么就更新一下执行信息,没有就启动。 - // 4. repair过程中下载的文件先放到delta文件夹下,然后再放到data目录下。 - Status repairTablet(RepairTabletTask task) - - // 这个API只是这里假定的,这个应该跟BE现有的Load API 合成一个 - Status loadData(dataUrl, transactionid) - // 当BE 收到这个消息时,读取本地staging目录下的文件的后缀名.stage3的文件,如果tableid 和 transactionid匹配,那么将文件,重命名,逻辑是: - // 1. 首先在 delta 目录下建立硬链指向staing目录下的文件。 - // 2. 判断一下data目录中version-1的文件是否存在,如果存在并且当前tablet是正常状态时,那么在data 目录中也建立硬链。 - // 3. 把staging目录下的文件删除。 - // 如果BE 修改version时发生问题,那么把异常的tablet返回给FE - void publishVersion(list <tuple<tableId, transactionid, version>>) - } - -# 相关流程和机制 -## 导入流程 - -由于这一阶段只针对现有的导入流程优化,所以我们从loading阶段开始描述我们的方案, loading任务之前的阶段保持不变。 - -- 当Extract和Transform阶段执行完毕后,FE会收到通知,那么FE开始执行Loading任务,此时FE充当了Coordinator的角色,FE调用beginTransaction API在事务表中注册这个导入任务。 -- FE 调用BE的loadData API 向BE发送loading 任务,在消息中要附加一个transactionid,表明这是一个新式的导入任务。【这里会不会成为瓶颈?但是以后实时导入不会有这个瓶颈,因为实时导入中load的通知是由coordinator通知的,是分散化的】对于通知不成功的,那么就一直轮训通知即可,实在是通知不成功,那么就把副本设置为异常状态即可。 -- BE 收到loading任务后执行loading任务,不论是小批量还是DPP,都从目标源(对于DPP来说就是HDFS,对于小批量来说就是那个BE)上下载文件, 但是要把数据放到 staging 文件夹下,文件的后缀名是.stage1。 -- 当BE执行完loading任务时,把文件的后缀名改成 .stage2 -- BE 启动一个定时服务,扫描 staging 文件夹下的.stage2 的文件,如果涉及到一个transaction的所有的tablet的导入都完成了,那么BE调用 finishLoad API向FE汇报load完成信息, FE在内存中记录这个导入完成情况,注意这里没有持久化。 -- BE 收到返回后,在内存中记录这个文件已经汇报过了,下次汇报的时候不再汇报,这样就能够避免每次都发送所有的导入完成的列表给FE了,但是这引入了一个新的问题,就是FE仅仅把状态保存在内存中,当FE重启或者FE迁移了内存状态就不在了,这个在FE故障处理章节介绍一下。 -- FE 启动一个定时任务扫描所有的导入任务,如果一个transaction涉及到的所有的tablet都导入完成,或者连续的quorum完成时,那么FE 调用commitTransaction API 来完成事务,这一步FE要判断一个BE是否是导入不成功了,因为commitTransaction那个API中会根据状态设置tablet的状态。注意这里没有采用触发式的,是为了降低FE的元数据修改次数。 - -此时导入流程完毕,但是BE端并不知道transactionid和version的对应关系,数据也不可以查询。后续由version 通知机制完成BE端从transactionid到version的变更。 - -## version 通知机制 - -这个服务只运行在FE Leader上,不是Leader不运行。 - -- 遍历transaction_state 表,选择可以通知的version, 选择的方式是: - - 事务的状态为COMMITTED - - 事务的tableidlist 中每个table的version都等于 table_version 表中的visible_version + 1 - - 还要检查一下副本的状态是不是健康,如果有不健康的,那么也直接把load任务设置为失败 -- 将遍历后获取的结果组装成list <tuple<tableId, transactionid, version>>格式,调用BE的publishVersoin API通知BE。 -- BE 收到通知后, 根据tableId和transactionid从staging目录下找导入文件,完成文件名的变更,具体看API的描述,如果遇到异常那么尝试3次,如果还有错误那么跳过,依赖补救措施来搞定。 -- 如果所有的BE 都返回成功,那么修改transaction_state 中的transaction修改为VISIBLE,同时把table_version 表中的visible_version 修改为对应的tableidlist 中记录的table的version。 -- 如果只有部分BE返回成功,那么把那些没返回成功的BE的状态 - -## version通知机制的保证/补救措施 - -在某些情况下可能publishVersion不能更新tablet的transactionid到version,比如我们假设一个现象,BE在接收publishVersion时,一个磁盘掉了,然后又回来了,那么这个磁盘上的tablet的version实际是没有变的,但是FE会误认为所有的tablet的version都更新成功了。 - -所以在这里我们引入一个强制保证的机制: - -- BE 定期的跟FE比对本地所有tablet的version,如果本地的version < FE中记录的visible version,那么判定本地是有问题的。 -- BE 扫描本地tablet的staging目录,从中获得所有的transactionid,从FE中获取对应的version,然后更新本地的version。 - -另外,当BE收到publishVersion做变更时,BE要检查下tablet对应的version-1数据是否存在,如果不存在,说明上次的publish并没有成功,可能有什么意外情况没考虑到,那么BE也需要从FE中主动拉一下transactionid对应的version。 - - -## BE 启动过程 - -- BE 启动时要完成自检的过程,跟FE通信汇报本地的tablet信息,如果一个磁盘坏了,那么需要向FE汇报,FE需要将这个tablet标记为异常。 - -## Compaction 逻辑 - -现在BE上一个tablet目录下的data目录里的数据 -保持现有的compaction逻辑不变,因为version的通知是顺序的,所以最后一个version仍然是未决version,仍然不compaction。 - -## delta目录的清空逻辑 - -delta目录下的数据单纯是为了副本修复服务的,所以我们这里采取简单的定时的策略,目前考虑半小时自动删除,另外如果磁盘空间不足的话,可以考虑删除。 - -## 副本故障处理 - -副本修复的触发机制和类型: - -- 当一个tablet的BE从故障中恢复汇报本地的tablet给FE时,FE决定这个tablet仍然由这个BE来存储时,FE直接生成RepairTabletTask。 -- 当Coordinator调用commitTransaction的API时,检测到一个tablet没有完成不完整时: - - 如果这个BE在元数据中标记为正常,那么此时立即生成一个RepairTabletTask。 - - 如果这个BE在元数据中标记为异常了,那么不执行副本修复。 -- FE定时检查副本的状态,如果一个副本处于故障状态超过一定的时间(比如15min,这个时间要参考两个机器同时宕机的概率),那么立即生成RepairTabletTask。 -- relbance 过程,当一个tablet从A机器迁移到B机器时,直接把B机器当做一个故障的副本,生成一个RepairTabletTask。 -- 副本数目调大,新加的副本认为是异常副本,生成一个RepaireTabletTask。 - -注意: 生成RepairTabletTask的同时把故障的tablet标记为clone状态,对于未来的loading任务也要向这个节点发送,version通知也要发送,只是在计算quorum的时候不计算。 - -RepairReplicaTask的执行流程: - -- FE 调用BE的repairTablet API,将RepairTabletTask发送给对应的BE。 -- BE 执行RepairTabletTask。 -- 当BE下载完后,跟data目录下的内容比对一下,如果完全对了,那么把数据在data目录下建立硬链,给FE汇报成功消息。 - -# 一些异常处理的思路 -## BE 的一个磁盘掉了,然后重启了,收到publishVersion 消息后如何处理? -此时我们不做处理,仍然认为version通知成功了。 其实现在的BE也是没处理的,这相当于BE给FE汇报导入成功后,BE自己的磁盘又掉了的情况,我们不做处理,等BE的定时汇报由FE发现异常,标记tablet为异常。 - -## FE 宕机重启 -- 内存中保存的导入完成情况丢失: 此时FE 成为Leader时,需要向所有的BE发送一个invalid消息,告诉BE本地的内存状态无效,从硬盘上汇报所有的导入进度信息给FE。 -- RepairTabletTask丢失:在创建RepairTabletTask时把tablet设置为clone状态,此时可以遍历CLONE状态的tablet生成RepairTabletTask,让BE重新做即可,这里的sourceIP可能发生变化,但是我们不管了。 - -# 遗留的导入的改造方式 -## 导入的改进方式 -- FE Leader 重启时根据元数据中已经loading finished的任务的version,将这个version的最大值根据table写入table_version 的visible_version 字段中。 -- FE Leader 根据正在loading的任务【尚未完成导入,ET阶段已经结束,Loading阶段未结束】,获取这些任务的version,获取最大的,把version写入table_version 的next_version 字段中。 -- 对于已经开始执行loading阶段的导入,FE Leader继续用轮训的方式,让他们执行loading任务,在loading完成后,要增加修改table对应的visible_version的逻辑。 这里可能不仅仅是quorum导入完成,要等尽量长的时间,让所有的副本都导入成功才可以。 如果副本导入没成功,那么标记为故障。 【或许这一步我们可以直接cancel掉之前所有的导入】 -- FE Leader上把现在追副本的任务停止。 - -# 其他一些考虑 -- 处于loading 阶段的任务不能太多,否则BE端压力太大,过去version lock能达到这个效果,现在没有version lock了,需要增加一个限制。 我们是不是有这个机制了呢? \ No newline at end of file diff --git a/docs/design/show_current_sql.md b/docs/design/show_current_sql.md deleted file mode 100644 index 540cec2c..00000000 --- a/docs/design/show_current_sql.md +++ /dev/null @@ -1,65 +0,0 @@ -# 查看正在执行的SQL功能设计 - -当前Palo在线上出现压力过大查询卡住时,难以定位当时是哪些查询导致集群压力过大,为了能在发生问题时,快速的定位并找到集群正在执行的大查询,所以需要提供查看正在执行的SQL的功能。 - -## 用户接口 - -### 1、查看正在执行的查询 - - show proc "/current_queries" - - - -### 2、查看某个查询的具体SQL - - show proc "/current_queries/{query_id}" - - - -### 3、查看某个查询的FRAGMENT的执行情况 - - show proc "/current_queries/{query_id}/fragments" - - - -FRAGMENT ID与EXPLAIN里面的FRAGMENT ID是一一对应的,方便查看查询卡到哪一步。 - -### 4、查看所有BE上正在执行的INSTANCE - - show proc "/current_backend_instances" - - - - - -## 详细设计 - -为了可以实现上述的功能,牵扯了三个方面: - -* 查询信息的统计 -* BE执行状态的引入 -* 查询的展示 - -查询的展示不做过多的说明,即通过实现ProcDirInterface接口然后注册到ProcService来实现4种接口的schema展示。 - -### 查询信息的统计 - -为了统计查询,在SQL下发执行时取到SQL、FRAGMENT的所有信息,这里需要有一个全局的查询信息统计的结构,重写了一下现有的QeProcessor作为一个接口: - - public interface QeProcessor { - - TReportExecStatusResult reportExecStatus(TReportExecStatusParams params); - - void registerQuery(TUniqueId queryId, Coordinator coord) throws InternalException; - - void registerQuery(TUniqueId queryId, QeProcessorImpl.QueryInfo info) throws InternalException; - - void unregisterQuery(TUniqueId queryId); - - Map<String, QueryStatisticsItem> getQueryStatistics(); - } - -通过registerQuery来注册查询,unregisterQuery来反注册查询,然后proc下的各个实现都调用getQueryStatistics来获取当前查询的信息,对于接口3需要查询SQL的INSTANCE在BE上的执行情况,所以对于BE的执行这里引入了两种状态。 - -* running : 正在进行计算 -* wait:等待数据的状态,对于EXCHANGE以及SCANNODE来讲读取数据的时候需要设置为此状态。 diff --git a/docs/design/show_load_warnings.md b/docs/design/show_load_warnings.md deleted file mode 100644 index d5fc2cde..00000000 --- a/docs/design/show_load_warnings.md +++ /dev/null @@ -1,182 +0,0 @@ -[TOC] - -# Show Load Warnings - -## 背景 - -在palo3.1以后,Palo存在三种数据导入方式: -+ 通过hadoop load批量导入; -+ 通过http进行小批量导入; -+ 使用broker进行数据导入; - -对于用户来说,查看导入时的错误信息非常必要。根据这些错误信息,用户进行数据的整理和修改。 - -目前Palo支持的导入方式和对应的错误详情查看方式如下表: - -| 导入方式 | 内网 | 公有云 | -| :---------: | :--------------------------------: | :------: | -| Hadoop | HTTP,访问be | 无 | -| 小批量导入 | HTTP,访问map-reduce作业的error_log | 无 | -| Pull-Load | ---------- | ------ | - -这里有两个问题,一个是**没有提供统一的导入错误查看接口**。另一个问题,也是最严重的,**对于公有云用户,目前无法自行查看导入的错误详细信息。** - -## 需求 - -无论是通过何种方式进行导入的数据,用户**有统一的接口进行查看错误信息,并且在内网和外网都可以进行查看。** - -用户查看错误信息都是以“数据行”为单位,所以信息展示时也以“数据行”为单位。 - -但是对于一个导入任务,不需要展示其所有的数据数据行,只需要展示部分即可。总的正确行数和错误行数仍然在SHOW LOAD的EtlInfo中查看。 - -## 方案概述 - -### 1. 对外复用现有的mysql协议。 - -因为我们对用户提供的是mysql协议接口,所以为了保证内网和外网用户都能访问,使用mysql协议输出即可。这既复用了现有的协议,又不需要额外部署服务。 - -### 2. 内部将这些错误信息输出到统一的“介质”。 - -在一个导入任务中,可能会有多台机器进行ETL。这些机器可能是BE(小批量导入时),也可能是hadoop的机器(map-reduce任务的mapper),每台机器保存自己处理数据时遇到的错误,即**这些错误信息分布在多台机器上。** - -如果在用户请求时再进行收集和展示,需要向多台机器发送请求,效率不高。 - -另外,M-R作业的error-log,还需要解析日志文件,然后处理后才能用于展示,比较繁琐。 - -为了解决这个问题,我们在错误信息生成时,就进行统一收集,而不再保存在本地。 -即当每台机器在执行ETL时发现原始数据有误时,都将相应的错误信息数据到相同的地方。所以这个“统一的介质”**是一个服务**,收集其他机器推送的数据。 - -另一方面,这个“统一的介质”**要能够进行高效的查询**,从而当用户查询时能够快速的响应。 - -## 详细设计 - -### 1. mysql接口设计 - -``` sql -SHOW LOAD WARNINGS -[FROM DB_NAME] -[ - WHERE - [LABEL = "your_label"] - [LOAD_JOB_ID = your_job_id] -] -[LIMIT limit]; -``` - -> 注:这里使用Warning,而没有使用Error。因为当用户允许一定的filter时,即使有错误的数据行,但导入任务并不一定会失败,这符合Warning的语义,而Error的感觉是一旦出现error导入作业一定失败。 - -说明: -1) 如果不指定 db_name,使用当前默认db -2) 如果使用 LABEL = ,则精确匹配指定的 label。当导入失败重试时,同一个label可以对应多个job,默认只显示最新的一个job。 -3) 如果使用 LOAD_JOB_ID = ,则精确匹配指定的 job -4) 如果指定了 LIMIT,则显示 limit 条匹配记录。否则全部显示。(后面可以考虑,即使不设置limit,也最多展示一定条数,比如100行) - -### 2. 结果显示 - -| JobId | Label| ErrorMsgDetail | -| :---: | :--: | :------------: | -| *** | **** | ********** | - -每行都对应原始数据中的一行。 - -### 3. 设置mysql服务地址 - -格式: - -``` sql -ALTER SYSTEM SET LOAD_ERROR_URL= "mysql://user:password@host:port[/database[/table]]" -``` - -如果以后不再使用mysql,只需要将"mysql"修改对应的协议栈即可。 - -如果不指定database,那么默认使用$cluster_id作为数据库名; -如果不指定table,那么默认使用"load_errors"作为表名。 - -### 4. 底层引擎选择 - -如上所述,“介质”是一个服务,还需要能够高效查询。 - -为了实现简单,可以Mysql或者ES,我们这次使用Mysql。 - -如果以后Mysql性能成为瓶颈,可以再考虑ES,或者自己实现一个简单存储服务,比如利用rocksdb,或者自己封装一个HashMap。 - -部署Mysql时,要保证能够同时被M-R的mapper和FE&BE都可以访问。公司内网的palo集群,就在内网部署mysql; 公有云的palo集群,就在公有云可访问的网段部署; - -### 5. 表设计 - -```sql -create table load_errors ( - job_id BIGINT NOT NULL, - error_msg VARCHAR(100) NOT NULL, - INDEX(job_id) -) -``` - -为了便于快速查询,**在job_id列上建立了index。** - -将来如果需要,可以再增加一个type列,表示错误类型。(目前小批量导入,没有严格的统计区分错误类型。) - -### 6. 数据量预估 - -目前线上nmg集群导入量最大,7天的导入任务数量为13.5万,平均每天不到2w。 - -假设每个任务平均最大输出为100条,那么每天错误信息的条数约为200W,Mysql应该不会成为瓶颈。 - -### 7. 插入方式 - -```sql -INSERT INTO load_errors (job_id, error_msg) -VALUES( **, "" ); -``` - -### 8. 查询方式 - -```sql -SELECT job_id, error_msg -FROM load_errors -WHERE job_id in (xx, xx); -``` - -### 9. TTL(历史数据清理) - -为了简单起见,我们按时清理历史数据。执行和删除label相同的逻辑即可,保留7天。 - -这里并不需要和删除label耦合,即如果错误信息的条数增长很快,那么可以将保留时间降低,从而提高效率。 - -注意一点:删除时考虑到Mysql的性能,不进行批量删除,而是利用jobid的递增特性,每次删除100条,发多次。 - -### 10. 可靠性保证 - -这里保留的错误信息只需要提供**尽力而为**的原则。即便是数据丢失,是不会对系统运行造成不良影响的。所以我们暂时不需要部署Mysql Slave。 - -如果Mysql服务宕机,我们可以很快的搭建一个新的Mysql继续提供服务。 - -为了能够保证Mysql切换时不需要重新修改线上配置,Mysql的地址使用BGW。这样切换机器时,只需要修改域名映射即可。 - -另外,在向Mysql写入时,为了不阻塞导入,也是采用尽力而为的写法,如果写入mysql失败,那么就放弃写入,导入作业可以继续进行。 - -### 11. 多集群共用 - -目前在公有云,每个用户都会创建自己的Palo集群,不可能为每个集群单独部署mysql服务。另外,公司内网也有多个palo集群,可以共用一个Mysql。 - -我们部署的mysql是公用的,这需要一定的隔离机制。 - -在Palo中,公司内网用户通过DB进行隔离;公有云的用户通过集群进行隔离。 - -为了统一,在导入错误信息的Mysql服务中,可以**通过DB来区分不同集群**,即在Mysql中每个集群对应一个DB。 - -**DB的命名为:{$cluster_id}** - -说明:因为cluster_id是一个随机的64位整数,理论上不同的集群上是可能有相同的,但是概率极低。另外,即便是重复了,结果就是这两个集群会共用一个集群。只有在jobid也重复的情况下,才会出现显示内容混乱,这个概率就更低了。 - -这样,在DPP,FE,BE中,都只需要感知 cluster_id 就可以定位到对应的DB。 - -## 实现 - -主要分以下几个方面: -1. FE:Java, 实现解析用户"show load error"命令,拼装并查询Mysql,然后组装成最后结果返回用户。load_error_url信息要持久化。 -2. BE:C++, 封装写入接口,共小批量导入和Pull_Load调用。 -3. DPP:Python, 封装写入接口,mapper在输出错误时调用。 - - - diff --git a/docs/help/Contents/Administration/admin_show_stmt.md b/docs/help/Contents/Administration/admin_show_stmt.md new file mode 100644 index 00000000..1cde7683 --- /dev/null +++ b/docs/help/Contents/Administration/admin_show_stmt.md @@ -0,0 +1,64 @@ +# ADMIN SHOW REPLICA STATUS +## description + + 该语句用于展示一个表或分区的副本状态信息 + + 语法: + + ADMIN SHOW REPLICA STATUS FROM [db_name.]tbl_name [PARTITION (p1, ...)] + [where_clause]; + + where_clause: + WHERE STATUS [!]= "replica_status" + + replica_status: + OK: replica 处于监控状态 + DEAD: replica 所在 Backend 不可用 + VERSION_ERROR: replica 数据版本有缺失 + MISSING: replica 不存在 + +## example + + 1. 查看表全部的副本状态 + + ADMIN SHOW REPLICA STATUS FROM db1.tbl1; + + 2. 查看表某个分区状态为 VERSION_ERROR 的副本 + + ADMIN SHOW REPLICA STATUS FROM tbl1 PARTITION (p1, p2) + WHERE STATUS = "VERSION_ERROR"; + + 3. 查看表所有状态不健康的副本 + + ADMIN SHOW REPLICA STATUS FROM tbl1 + WHERE STATUS != "OK"; + +## keyword + ADMIN,SHOW,REPLICA,STATUS + +# ADMIN SHOW REPLICA DISTRIBUTION +## description + + 该语句用于展示一个表或分区副本分布状态 + + 语法: + + ADMIN SHOW REPLICA DISTRIBUTION FROM [db_name.]tbl_name [PARTITION (p1, ...)]; + + 说明: + + 结果中的 Graph 列以图形的形式展示副本分布比例 + +## example + + 1. 查看表的副本分布 + + ADMIN SHOW REPLICA DISTRIBUTION FROM tbl1; + + 1. 查看表的分区的副本分布 + + ADMIN SHOW REPLICA DISTRIBUTION FROM db1.tbl1 PARTITION(p1, p2); + +## keyword + ADMIN,SHOW,REPLICA,DISTRIBUTION + diff --git a/docs/help/Contents/Data Manipulation/manipulation_stmt.md b/docs/help/Contents/Data Manipulation/manipulation_stmt.md index dee0ceda..95650621 100644 --- a/docs/help/Contents/Data Manipulation/manipulation_stmt.md +++ b/docs/help/Contents/Data Manipulation/manipulation_stmt.md @@ -905,3 +905,17 @@ ## keyword SHOW, SNAPSHOT + +# RESTORE TABLET +## description + + 该功能用于恢复trash目录中被误删的tablet数据。 + + 说明:这个功能暂时只在be服务中提供一个http接口。如果要使用, + 需要向要进行数据恢复的那台be机器的http端口发送restore tablet api请求。api格式如下: + METHOD: POST + URI: http://be_host:be_http_port/api/restore_tablet?tablet_id=xxx&schema_hash=xxx + +## example + + curl -X POST "http://hostname:8088/api/restore_tablet?tablet_id=123456&schema_hash=1111111" diff --git a/docs/help/Contents/Data Manipulation/streaming.md b/docs/help/Contents/Data Manipulation/streaming.md new file mode 100644 index 00000000..4e31a1db --- /dev/null +++ b/docs/help/Contents/Data Manipulation/streaming.md @@ -0,0 +1,150 @@ +# STREAM LOAD +## description + NAME: + stream-load: load data to table in streaming + + SYNOPSIS + curl --location-trusted -u user:passwd [-H ""...] -T data.file -XPUT http://fe_host:http_port/api/{db}/{table}/_stream_load + + DESCRIPTION + 该语句用于向指定的 table 导入数据,与普通Load区别是,这种导入方式是同步导入。 + 这种导入方式仍然能够保证一批导入任务的原子性,要么全部数据导入成功,要么全部失败。 + 该操作会同时更新和此 base table 相关的 rollup table 的数据。 + 这是一个同步操作,整个数据导入工作完成后返回给用户导入结果。 + 当前支持HTTP chunked与非chunked上传两种方式,对于非chunked方式,必须要有Content-Length来标示上传内容长度,这样能够保证数据的完整性。 + 另外,用户最好设置Expect Header字段内容100-continue,这样可以在某些出错场景下避免不必要的数据传输。 + + OPTIONS + 用户可以通过HTTP的Header部分来传入导入参数 + + label: 一次导入的标签,相同标签的数据无法多次导入。用户可以通过指定Label的方式来避免一份数据重复导入的问题。 + 当前Palo内部保留30分钟内最近成功的label。 + + column_separator:用于指定导入文件中的列分隔符,默认为\t。如果是不可见字符,则需要加\x作为前缀,使用十六进制来表示分隔符。 + 如hive文件的分隔符\x01,需要指定为-H "column_separator:\x01" + + columns:用于指定导入文件中的列和 table 中的列的对应关系。如果源文件中的列正好对应表中的内容,那么是不需要指定这个字段的内容的。 + 如果源文件与表schema不对应,那么需要这个字段进行一些数据转换。这里有两种形式column,一种是直接对应导入文件中的字段,直接使用字段名表示; + 一种是衍生列,语法为 `column_name` = expression。举几个例子帮助理解。 + 例1: 表中有3个列“c1, c2, c3”,源文件中的三个列一次对应的是"c3,c2,c1"; 那么需要指定-H "columns: c3, c2, c1" + 例2: 表中有3个列“c1, c2, c3", 源文件中前三列依次对应,但是有多余1列;那么需要指定-H "columns: c1, c2, c3, xxx"; + 最后一个列随意指定个名称占位即可 + 例3: 表中有3个列“year, month, day"三个列,源文件中只有一个时间列,为”2018-06-01 01:02:03“格式; + 那么可以指定-H "columns: col, year = year(col), month=mont(col), day=day(col)"完成导入 + + where: 用于抽取部分数据。用户如果有需要将不需要的数据过滤掉,那么可以通过设定这个选项来达到。 + 例1: 只导入大于k1列等于20180601的数据,那么可以在导入时候指定-H "where: k1 = 20180601" + + max_filter_ratio:最大容忍可过滤(数据不规范等原因)的数据比例。默认零容忍。 + partitions: 用于指定这次导入所设计的partition。如果用户能够确定数据对应的partition,推荐指定该项。不满足这些分区的数据将被过滤掉。 + 比如指定导入到p1, p2分区,-H "partitions: p1, p2" + + RETURN VALUES + 导入完成后,会以Json格式返回这次导入的相关内容。当前包括一下字段 + Status: 导入最后的状态。 + Success:表示导入成功,数据已经可见; + Publish Timeout:表述导入作业已经成功Commit,但是由于某种原因并不能立即可见。用户可以视作已经成功不必重试导入 + Label Already Exists: 表明该Label已经被其他作业占用,可能是导入成功,也可能是正在导入。 + 用户需要通过get label state命令来确定后续的操作 + 其他:此次导入失败,用户可以指定Label重试此次作业 + Message: 导入状态详细的说明。失败时会返回具体的失败原因。 + NumberLoadedRows: 此次导入的数据行数,只有在Success时有效 + NumberFilteredRows: 此次导入过滤掉的行数 + LoadBytes: 此次导入的源文件数据量大小 + LoadTimeMs: 此次导入所用的时间 + ErrorURL: 被过滤数据的具体内容,仅保留前1000条 + + ERRORS + +## example + + 1. 将本地文件'testData'中的数据导入到数据库'testDb'中'testTbl'的表,使用Label用于去重 + curl --location-trusted -u root -H "lable:123" -T testData http://host:port/api/testDb/testTbl/_stream_load + + 2. 将本地文件'testData'中的数据导入到数据库'testDb'中'testTbl'的表,使用Label用于去重, 并且只导入k1等于20180601的数据 + curl --location-trusted -u root -H "lable:123" -H "where: k1=20180601" -T testData http://host:port/api/testDb/testTbl/_stream_load + + 3. 将本地文件'testData'中的数据导入到数据库'testDb'中'testTbl'的表, 允许20%的错误率(用户是defalut_cluster中的) + curl --location-trusted -u root -H "lable:123" -H "max_filter_ratio:0.2" -T testData http://host:port/api/testDb/testTbl/_stream_load + + 4. 将本地文件'testData'中的数据导入到数据库'testDb'中'testTbl'的表, 允许20%的错误率,并且指定文件的列名(用户是defalut_cluster中的) + curl --location-trusted -u root -H "lable:123" -H "max_filter_ratio:0.2" -H "columns: k2, k1, v1" -T testData http://host:port/api/testDb/testTbl/_stream_load + + 5. 将本地文件'testData'中的数据导入到数据库'testDb'中'testTbl'的表中的p1, p2分区, 允许20%的错误率。 + curl --location-trusted -u root -H "lable:123" -H "max_filter_ratio:0.2" -H "partitions: p1, p2" -T testData http://host:port/api/testDb/testTbl/_stream_load + + 6. 使用streaming方式导入(用户是defalut_cluster中的) + seq 1 10 | awk '{OFS="\t"}{print $1, $1 * 10}' | curl --location-trusted -u root -T - http://host:port/api/testDb/testTbl/_stream_load + + 7. 导入含有HLL列的表,可以是表中的列或者数据中的列用于生成HLL列 + curl --location-trusted -u root -H "columns: k1, k2, v1=hll_hash(k1)" -T testData http://host:port/api/testDb/testTbl/_stream_load + +## keyword + STREAM,LOAD + +# GET LABEL STATE +## description + NAME: + get_label_state: get label's state + + SYNOPSIS + curl -u user:passwd http://host:port/api/{db}/{label}/_state + + DESCRIPTION + 该命令用于查看一个Label对应的事务状态 + + RETURN VALUES + 执行完毕后,会以Json格式返回这次导入的相关内容。当前包括一下字段 + Label:本次导入的 label,如果没有指定,则为一个 uuid。 + Status:此命令是否成功执行,Success表示成功执行 + Message: 具体的执行信息 + State: 只有在Status为Success时才有意义 + UNKNOWN: 没有找到对应的Label + PREPARE: 对应的事务已经prepare,但尚未提交 + COMMITTED: 事务已经提交,不能被cancel + VISIBLE: 事务提交,并且数据可见,不能被cancel + ABORTED: 事务已经被ROLLBACK,导入已经失败。 + + ERRORS + +## example + + 1. 获得testDb, testLabel的状态 + curl -u root http://host:port/api/testDb/testLabel/_state + +## keyword + GET, LABEL, STATE + +# CANCEL LABEL +## description + NAME: + cancel_label: cancel a transaction with label + + SYNOPSIS + curl -u user:passwd -XPOST http://host:port/api/{db}/{label}/_cancel + + DESCRIPTION + 该命令用于cancel一个指定Label对应的事务,事务在Prepare阶段能够被成功cancel + + RETURN VALUES + 执行完成后,会以Json格式返回这次导入的相关内容。当前包括一下字段 + Status: 是否成功cancel + Success: 成功cancel事务 + 其他: cancel失败 + Message: 具体的失败信息 + + ERRORS + +## example + + 1. cancel testDb, testLabel的作业 + curl -u root -XPOST http://host:port/api/testDb/testLabel/_cancel + +## keyword + CANCEL,LABEL + + + + + + diff --git a/docs/help/Contents/Utility/util_stmt.md b/docs/help/Contents/Utility/util_stmt.md index 0df0f2d9..2fbafe82 100644 --- a/docs/help/Contents/Utility/util_stmt.md +++ b/docs/help/Contents/Utility/util_stmt.md @@ -1,13 +1,13 @@ -# DESCRIBE -## description - 该语句用于展示指定 table 的 schema 信息 - 语法: - DESC[RIBE] [db_name.]table_name [ALL]; - - 说明: - 如果指定 ALL,则显示该 table 的所有 index 的 schema - -## example - -## keyword +# DESCRIBE +## description + 该语句用于展示指定 table 的 schema 信息 + 语法: + DESC[RIBE] [db_name.]table_name [ALL]; + + 说明: + 如果指定 ALL,则显示该 table 的所有 index 的 schema + +## example + +## keyword DESCRIBE,DESC \ No newline at end of file diff --git a/gensrc/script/Makefile b/gensrc/script/Makefile index 0686bc2f..11725bbf 100644 --- a/gensrc/script/Makefile +++ b/gensrc/script/Makefile @@ -50,14 +50,14 @@ GEN_OPCODE_OUTPUT = ${BUILD_DIR}/thrift/Opcodes.thrift \ ${BUILD_DIR}/java/org/apache/doris/opcode/FunctionRegistry.java \ ${BUILD_DIR}/java/org/apache/doris/opcode/FunctionOperator.java -${GEN_OPCODE_OUTPUT}: palo_functions.py ${GEN_FUNC_OUTPUT} ${GEN_VEC_FUNC_OUTPUT} | ${BUILD_DIR}/python +${GEN_OPCODE_OUTPUT}: doris_functions.py ${GEN_FUNC_OUTPUT} ${GEN_VEC_FUNC_OUTPUT} | ${BUILD_DIR}/python gen_opcode: ${GEN_OPCODE_OUTPUT} .PHONY: gen_opcode # generate GEN_BUILTINS_OUTPUT = ${BUILD_DIR}/java/org/apache/doris/builtins/ScalarBuiltins.java -${GEN_BUILTINS_OUTPUT}: palo_builtins_functions.py gen_builtins_functions.py +${GEN_BUILTINS_OUTPUT}: doris_builtins_functions.py gen_builtins_functions.py cd ${BUILD_DIR}/python && ${PYTHON} ${CURDIR}/gen_builtins_functions.py gen_builtins: ${GEN_BUILTINS_OUTPUT} .PHONY: gen_builtins diff --git a/gensrc/script/palo_builtins_functions.py b/gensrc/script/doris_builtins_functions.py similarity index 100% rename from gensrc/script/palo_builtins_functions.py rename to gensrc/script/doris_builtins_functions.py diff --git a/gensrc/script/palo_functions.py b/gensrc/script/doris_functions.py similarity index 100% rename from gensrc/script/palo_functions.py rename to gensrc/script/doris_functions.py diff --git a/gensrc/script/gen_builtins_functions.py b/gensrc/script/gen_builtins_functions.py index f13cdc23..6cef755a 100755 --- a/gensrc/script/gen_builtins_functions.py +++ b/gensrc/script/gen_builtins_functions.py @@ -1,11 +1,11 @@ """ -This module is palo builtin functions +This module is doris builtin functions """ import sys import os from string import Template -import palo_builtins_functions +import doris_builtins_functions java_registry_preamble = '\ // Licensed to the Apache Software Foundation (ASF) under one \n\ @@ -27,7 +27,7 @@ // This is a generated file, DO NOT EDIT.\n\ // To add new functions, see the generator at\n\ // common/function-registry/gen_builtins_catalog.py or the function list at\n\ -// common/function-registry/palo_builtins_functions.py.\n\ +// common/function-registry/doris_builtins_functions.py.\n\ \n\ package org.apache.doris.builtins;\n\ \n\ @@ -42,7 +42,7 @@ }\n\ }\n' -FE_PATH = "../java/org.apache.doris/builtins/" +FE_PATH = "../java/org/apache/doris/builtins/" # This contains all the metadata to describe all the builtins. # Each meta data entry is itself a map to store all the meta data @@ -54,7 +54,7 @@ def add_function(fn_meta_data, user_visible): """add function """ assert 4 <= len(fn_meta_data) <= 6, \ - "Invalid function entry in palo_builtins_functions.py:\n\t" + repr(fn_meta_data) + "Invalid function entry in doris_builtins_functions.py:\n\t" + repr(fn_meta_data) entry = {} entry["sql_names"] = fn_meta_data[0] entry["ret_type"] = fn_meta_data[1] @@ -115,9 +115,9 @@ def generate_fe_registry_init(filename): java_registry_file.close() # Read the function metadata inputs -for function in palo_builtins_functions.visible_functions: +for function in doris_builtins_functions.visible_functions: add_function(function, True) -for function in palo_builtins_functions.invisible_functions: +for function in doris_builtins_functions.invisible_functions: add_function(function, False) if not os.path.exists(FE_PATH): diff --git a/gensrc/script/gen_opcodes.py b/gensrc/script/gen_opcodes.py index 24df9025..d0b2e987 100755 --- a/gensrc/script/gen_opcodes.py +++ b/gensrc/script/gen_opcodes.py @@ -23,7 +23,7 @@ # about type checking. # # This scripts pulls function metadata input from -# - src/common/function/palo_functions.py (manually maintained) +# - src/common/function/doris_functions.py (manually maintained) # - src/common/function/generated_functions.py (auto-generated metadata) # # This script will generate 4 outputs @@ -44,7 +44,7 @@ import os import string sys.path.append(os.getcwd()) -import palo_functions +import doris_functions import generated_functions import generated_vector_functions @@ -337,15 +337,15 @@ def generate_fe_registry_init(filename): java_registry_file.close() # Read the function metadata inputs -for function in palo_functions.functions: +for function in doris_functions.functions: if len(function) != 5: - print "Invalid function entry in palo_functions.py:\n\t" + repr(function) + print "Invalid function entry in doris_functions.py:\n\t" + repr(function) sys.exit(1) add_function(function, False) -for function in palo_functions.udf_functions: +for function in doris_functions.udf_functions: assert len(function) == 6, \ - "Invalid function entry in palo_functions.py:\n\t" + repr(function) + "Invalid function entry in doris_functions.py:\n\t" + repr(function) add_function(function, True) for function in generated_functions.functions: ---------------------------------------------------------------- This is an automated message from the Apache Git Service. To respond to the message, please log on GitHub and use the URL above to go to the specific comment. For queries about this service, please contact Infrastructure at: [email protected] With regards, Apache Git Services --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
