大数据的处理之数据的抽取
学习目标:
数据抽取的方式和实现方法
学习内容:
- 数据的抽取方式:全量抽取和增量抽取
2.数据加载的方法:全表删除插入方式,触发器方式
学习时间:
有sql 基础的话 ,6个小时
学习产出:
1.技术笔记一篇
2.练习题一套,包括答案 源代码
一 概念
(1)全量抽取
全量抽取类似于数据迁移或数据复制,它将数据源中的表的数据原封不动的从数据库中抽取出来,并转换成自己的ETL工具可以识别的格式。全量抽取比较简单。
(2)增量抽取
增量抽取只抽取自上次抽取以来数据库中要抽取的表中新增或修改的数据。在ETL使用过程中。增量抽取较全量抽取应用更广。如何捕获变化的数据是增量抽取的关键。对捕获方法一般有两点要求:准确性,能够将业务系统中的变化数据按一定的频率准确地捕获到;不能对业务系统造成太大的压力,影响现有业务。
二 数据加载的方法
1.全表删除插入方式是指每次抽取前先删除目标表数据,抽取时全新加载数据。该方式实际上将增量抽取等同于全量抽取。对于数据量不大,全量抽取的时间代价小于执行增量抽取的算法和条件代价时,可以采用该方式
2.触发器方式
触发器方式是普遍采取的一种增量抽取机制。该方式是根据抽取要求,在要被抽取的源表上建立插入、修改、删除3个触发器,每当源表中的数据发生变化,就被相 应的触发器将变化的数据写入一个增量日志表,ETL的增量抽取则是从增量日志表中而不是直接在源表中抽取数据,同时增量日志表中抽取过的数据要及时被标记 或删除。
3.对于分区表 只抽取一个分区的内容。
例子从数据下载平台,另外的数据库中提取数据到自己的数据库。
把用户testkz表PK_FBK_OPEN数据导入到目标库用户sys中PK_FBK_OPEN
1.在用户testkz,sys建立表格
create table PK_FBK_OPEN
(
obk_no VARCHAR2(14),
op_date NUMBER(6),
op_inst VARCHAR2(8),
paper_no VARCHAR2(18),
paper_type VARCHAR2(2)
)
插入数据
insert into PK_FBK_OPEN
select ‘130001’,‘201906’,‘121205’,‘13010519820605’,‘01’ from dual
insert into PK_FBK_OPEN
select ‘130002’,‘201908’,‘121206’,‘130105198207605’,‘02’ from dual;
Commit;
2.建立数据库链
CREATE DATABASE LINK db_testkz
CONNECT TO testkz IDENTIFIED BY haien
USING ‘orcl’;
3建立配置表
create table LOAD_TAB_CONF
(
id VARCHAR2(50) default sys_guid(),
tab_name VARCHAR2(30),
tab_owner VARCHAR2(30),
tab_src VARCHAR2(30),
db_link VARCHAR2(30),
is_full CHAR(1),
is_enable CHAR(1)
);
comment on column LOAD_TAB_CONF.tab_name
is ‘源表名’;
comment on column LOAD_TAB_CONF.tab_owner
is ‘源表所有者’;
comment on column LOAD_TAB_CONF.tab_src
is ‘数据来源’;
comment on column LOAD_TAB_CONF.db_link
is ‘DBlink名称’;
comment on column LOAD_TAB_CONF.is_full
is ‘是否是全量1-是;0-否;2-静态表’;
comment on column LOAD_TAB_CONF.is_enable
is ‘是否加载’;
insert into LOAD_TAB_CONF(id,tab_name,tab_owner,tab_src,db_link,is_full,is_enable)
select 000001,‘PK_FBK_OPEN’,‘testkz’,‘数据平台’,db_testkz,1,1 from dual;
commit;
3.写存储过程
create or replace procedure PRO_LOAD_INTO_LOCAL(rp_dt char --日期 6位
) is
v_sql_text varchar2(4000);
v_colname varchar2(30);
v_ymq varchar2(14);
v_type_log type_log default type_log(rp_dt,
'PRO_LOAD_INTO_LOCAL',
null,
null,
null,
null,
null);
begin
v_ymq := to_date(to_char(last_day(to_date(rp_dt, ‘yyyymm’)), ‘yyyymmdd’),
‘yyyymmdd’);
–写日志
v_sql_text := ‘select ‘‘插入一条日志’’ from dual’;
v_type_log.run_comm := ‘开始将数据抽取到本地’;
pro_execute_sql(v_sql_text, v_type_log);
commit;
v_type_log.run_comm := ‘正在将数据抽取到本地…’;
for day_tab in (select t.tab_name, t.tab_owner, t.db_link, is_full
from load_tab_conf t
where t.is_enable = ‘1’ ) loop
–全量数据
if (day_tab.is_full = ‘1’) then
–清理数据
v_sql_text := ‘truncate table ’ || day_tab.tab_name;
pro_execute_sql(v_sql_text, v_type_log);
commit;
–抽取数据
if (day_tab.tab_name = ‘T_CARD_DTL’) then --只取当月最后一日数据即可
v_sql_text := ‘insert /+ append parallel(t,4) / into ’ ||
day_tab.tab_name ||
’ t select /+ parallel(s,4)/ * from ’ ||
day_tab.tab_owner || ‘.’ || day_tab.tab_name || ‘@’ ||
day_tab.db_link || ’ s where SUMM_DATE=’’’ || v_ymq || ‘’‘’;
–dbms_output.put_line(v_sql_text);
pro_execute_sql(v_sql_text, v_type_log);
commit;
else
v_sql_text := 'insert /*+ append parallel(t,4) */ into ' ||
day_tab.tab_name ||
' t select /*+ parallel(s,4)*/ * from ' ||
day_tab.tab_owner || '.' || day_tab.tab_name || '@' ||
day_tab.db_link || ' s';
--dbms_output.put_line(v_sql_text);
pro_execute_sql(v_sql_text, v_type_log);
commit;
end if;
--增量数据
elsif (day_tab.is_full = '0') then
--清理数据
v_sql_text := 'alter table ' || day_tab.tab_name ||
' truncate partition P_' || rp_dt;
pro_execute_sql(v_sql_text, v_type_log);
commit;
select t.column_name
into v_colname
from all_part_key_columns t
where t.name = day_tab.tab_name;
--dbms_output.put_line(v_colname);
--抽取数据
v_sql_text := 'insert /*+ append parallel(t,4) */ into ' ||
day_tab.tab_name ||
' t select /*+ parallel(s,4)*/ * from ' ||
day_tab.tab_owner || '.' || day_tab.tab_name || '@' ||
day_tab.db_link || ' s where s.v_colname>=''' || rp_dt ||
'01000000'' and s.v_colname<''' || v_ymq || '''';
--dbms_output.put_line(v_sql_text);
pro_execute_sql(v_sql_text, v_type_log);
commit;
--静态数据(系统初始化时装入,只装入一次)
else
--抽取数据
v_sql_text := 'insert /*+ append parallel(t,4) */ into ' ||
day_tab.tab_name ||
' t select /*+ parallel(s,4)*/ * from ' ||
day_tab.tab_owner || '.' || day_tab.tab_name || '@' ||
day_tab.db_link || ' s ';
pro_execute_sql(v_sql_text, v_type_log);
commit;
end if;
--改变启用状态
v_sql_text := 'update load_tab_conf t set t.is_enable =''0'' where t.tab_name=''' ||
day_tab.tab_name || '''';
--dbms_output.put_line(v_sql_text);
pro_execute_sql(v_sql_text, v_type_log);
commit;
end loop;
–循环完毕,除去静态表改变启用状态为1
v_sql_text := ‘update load_tab_conf t set t.is_enable =’‘1’’ where t.is_full<>‘‘2’’';
–dbms_output.put_line(v_sql_text);
pro_execute_sql(v_sql_text, v_type_log);
commit;
–写日志
v_type_log.run_comm := ‘将数据抽取到本地完成!’;
v_sql_text := ‘select ‘‘插入一条日志’’ from dual’;
pro_execute_sql(v_sql_text, v_type_log);
commit;
exception
when others then
rollback;
–写日志
v_type_log.run_comm := ‘将数据抽取到本地过程出错!’;
v_sql_text := ‘select ‘‘插入一条日志’’ from dual’;
pro_execute_sql(v_sql_text, v_type_log);
raise;
end PRO_LOAD_INTO_LOCAL;
附录 相关存储过程
create or replace procedure pro_execute_sql(i_sql clob,i_type_log in out type_log) is
v_begin integer;
v_end integer;
v_log_id integer;
begin
execute immediate ‘alter session enable parallel dml’;
i_type_log.sql_text:=i_sql;
i_type_log.log_level:=‘INFO’;
v_begin := dbms_utility.get_time();
execute immediate i_sql;
v_end :=dbms_utility.get_time;
i_type_log.cmt_count := sql%rowcount;
commit;
i_type_log.sql_elaps_tm:=(v_end-v_begin)/100;
pro_run_log(i_type_log,v_log_id);
exception when others then
rollback;
i_type_log.log_level:=‘ERR’;
pro_run_log(i_type_log,v_log_id);
raise;
end pro_execute_sql;
create or replace procedure pro_run_log(i_type type_log,o_id out integer) is
v_sqlcode varchar2(10);
begin
v_sqlcode :=sqlcode;
select seq_sys_run_log.nextval into o_id from dual;
insert into sys_run_log(
id ,
run_date,
pro_name,
log_level,
run_comm,
run_tm,
error_cd,
error_info,
sql_text,
sql_elaps_tm,
cmt_count)
values(
o_id,
i_type.run_date,
i_type.pro_name,
i_type.log_level,
i_type.run_comm,
sysdate,
case when i_type.log_level=‘ERR’ then v_sqlcode else null end,
case when i_type.log_level=‘ERR’ then
dbms_utility.format_error_stack()||dbms_utility.format_error_backtrace() else null end,
i_type.sql_text,
i_type.sql_elaps_tm,
i_type.cmt_count);
commit;
end pro_run_log;
更多推荐
所有评论(0)