· 10 years ago · Feb 26, 2016, 01:33 PM
1Four Steps Strategy for Incremental Updates in Apache Hive on Hadoop
2
3How to update records in Apache Hive
4By Greg Phillips on July 15th, 2014
5Incremental Updates
6Hadoop and Hive are quickly evolving to outgrow previous limitations for integration and data access.
7On the near-term development roadmap, we expect to see Hive supporting full CRUD operations (Insert, Select, Update, Delete). As we wait for these advancements, there is still a need to work with the current options—OVERWRITE or APPEND— for Hive table integration.
8
9The OVERWRITE option requires moving the complete record set from source to Hadoop. While this approach may work for smaller data sets, it may be prohibitive at scale.
10
11The APPEND option can limit data movement only to new or updated records. As true Inserts and Updates are not yet available in Hive, we need to consider a process of preventing duplicate records as Updates are appended to the cumulative record set.
12
13In this blog, we will look at a four-step strategy for appending Updates and Inserts from delimited and RDBMS sources to existing Hive table definitions. While there are several options within the Hadoop platform for achieving this goal, our focus will be on a process that uses standard SQL within the Hive toolset.
14
15IMPORTANT: For this blog, we assume that each source table will have a unique single or multi-key identifier and that a “modified_date†field is maintained for each record – either defined as part of the original source table or added as part of the ingest process.
16
17Hive Table Definition Options: External, Local and View
18External Tables are the combination of Hive table definitions and HDFS managed folders and files. The table definition exists independent from the data, so that, if the table is dropped, the HDFS folders and files remain in their original state.
19
20Local Tables are Hive tables that are directly tied to the source data. The data is physically tied to the table definition and will be deleted if the table is dropped.
21
22Views, as with traditional RDBMS, are stored SQL queries that support the same READ interaction as HIVE tables, yet they do not store any data of their own. Instead, the data are stored and sourced from the HIVE tables referenced in the stored SQL query.
23
24The following process outlines a workflow that leverages all of the above in four steps:
25
26Screen Shot 2014-07-14 at 2.24.01 PM
27Ingest. Complete (base_table) table movement followed by Change (incremental_table) records only.
28Reconcile. Creating a Single View of Base + Change records (reconcile_view) to reflect the most up-to-date record set.
29Compact. Creating a Reporting table (reporting_table) from the reconciled view.
30Purge. Replacing the Base table with Reporting table contents and deleting any previously processed Change records before the next Data Ingestion cycle.
31The tables and views that will be a part of the Incremental Update Workflow are:
32
33base_table: A HIVE Local table that initially holds all records from the source system. After the initial processing cycle, it will maintain a copy of the most up-to-date synchronized record set from the source. At the end of each processing cycle, it is overwritten by the reporting_table (as explained in the Step 4: Purge).
34incremental_table: A HIVE External table that holds the incremental change records (INSERTS and UPDATES) from the source system. At the end of each processing cycle, it is cleared of content (as explained in the Step 4: Purge).
35reconcile_view: A HIVE View that combines and reduces the base_table and incremental_table content to show only the most up-to-date records. It is used to populate the reporting_table (as explained in Step 3: Compact).
36reporting_table: A HIVE Local table that holds the most up-to-date records for reporting purposes. It is also used to overwrite the base_table at the end of each processing run.
37Step 1: Ingest
38Depending on whether direct access is available to the RDBMS source system, you may opt for either a File Processing method (when no direct access is available) or RDBMS Processing (when database client access is available).
39
40Regardless of the ingest option, the processing workflow in this article requires:
41
42One-time, initial load to move all data from source table to HIVE.
43On-going, “Change Only†data loads from the source table to HIVE.
44Below, both File Processing and Database-direct (SQOOP) ingest will be discussed.
45
46File Processing
47
48For this blog, we assume that a file or set of files within a folder will have a delimited format and will have been generated from a relational system (i.e. records have unique keys or identifiers).
49
50Files will need to be moved into HDFS using standard ingest options:
51
52WebHDFS: Primarily used when integrating with applications, a Web URL provides an Upload end-point into a designated HDFS folder.
53NFS: Appears as a standard network drive and allows end-users to use standard Copy-Paste operations to move files from standard file systems into HDFS.
54Once the initial set of records are moved into HDFS, subsequent scheduled events can move files containing only new Inserts and Updates.
55
56RDBMS Processing
57
58SQOOP is the JDBC-based utility for integrating with traditional databases. A SQOOP Import allows for the movement of data into either HDFS (a delimited format can be defined as part of the Import definition) or directly into a Hive table.
59
60The entire source table can be moved into HDFS or Hive using the “–table†parameter
61
62sqoop import --connect jdbc:teradata://{host name or ip address}/Database=retail --connection-manager org.apache.sqoop.teradata.TeradataConnManager --username dbc --password dbc --table SOURCE_TBL --target-dir /user/hive/incremental_table -m 1
63
64After the initial import, subsequent imports can leverage SQOOP’s native support for “Incremental Import†by using the “check-columnâ€, “incremental†and “last-value†parameters.
65
66sqoop import --connect jdbc:teradata://{host name or ip address}/Database=retail --connection-manager org.apache.sqoop.teradata.TeradataConnManager --username dbc --password dbc --table SOURCE_TBL --target-dir /user/hive/incremental_table -m 1--check-column modified_date --incremental lastmodified --last-value {last_import_date}
67
68Alternately, you can leverage the “query†parameter, and have SQL select statements limit the import to new or changed records only.
69
70sqoop import --connect jdbc:teradata://{host name or ip address}/Database=retail --connection-manager org.apache.sqoop.teradata.TeradataConnManager --username dbc --password dbc --target-dir /user/hive/incremental_table -m 1 --query 'select * from SOURCE_TBL where modified_date > {last_import_date} AND $CONDITIONS’
71
72Note: For the initial load, substitute “base_table†for “incremental_tableâ€. For all subsequent loads, use “incremental_tableâ€.
73
74Step 2: Reconcile
75In order to support an on-going reconciliation between current records in HIVE and new change records, two tables should be defined: base_table and incremental_table
76
77base_table
78
79The example below shows DDL for the Hive table “base_table†that will include any delimited files located in HDFS under the ‘/user/hive/base_table’ directory. This table will house the initial, complete record load from the source system. After the first processing run, it will house the on-going, most up-to-date set of records from the source system:
80[code language=â€SQLâ€]
81CREATE TABLE base_table (
82id string,
83field1 string,
84field2 string,
85field3 string,
86field4 string,
87field5 string,
88modified_date string
89)
90ROW FORMAT DELIMITED
91FIELDS TERMINATED BY ‘,’
92LOCATION ‘/user/hive/base_table’;
93[/code]
94
95incremental_table
96
97The DDL below shows an external Hive table “incremental_table†that will include any delimited files with incremental change records, located in HDFS under the ‘/user/hive/incremental_append’ directory:
98
99[code language=â€SQLâ€]
100CREATE EXTERNAL TABLE incremental_table (
101id string,
102field1 string,
103field2 string,
104field3 string,
105field4 string,
106field5 string,
107modified_date string
108)
109ROW FORMAT DELIMITED
110FIELDS TERMINATED BY ‘,’
111LOCATION ‘/user/hive/incremental_table’;
112[/code]
113
114reconcile_view
115
116This view combines record sets from both the Base (base_table) and Change (incremental_table) tables and is reduced only to the most recent records for each unique “idâ€.
117It is defined as follows:
118[code language=â€SQLâ€]
119CREATE VIEW reconcile_view AS
120SELECT t1.* FROM
121(SELECT * FROM base_table
122UNION ALL
123SELECT * FROM incremental_table) t1
124JOIN
125(SELECT id, max(modified_date) max_modified FROM
126(SELECT * FROM base_table
127UNION ALL
128SELECT * FROM incremental_table) t2
129GROUP BY id) s
130ON t1.id = s.id AND t1.modified_date = s.max_modified;
131[/code]
132
133Example
134
135The sample data below represents the UNION of both the base_table and incremental_table. Note, there are new updates for “id†values 1 and 2, which are found as the last two records in the table. The record for “id†3 remains unchanged.
136
137Screen Shot 2014-07-14 at 5.57.32 PM
138
139The reconcile_view should only show one record for each unique “idâ€, based on the latest “modified_date†field value.
140
141The resulting query from “select * from reconcile_view†shows only three records, based on both unique “id†and “modified_dateâ€
142
143Screen Shot 2014-07-14 at 5.57.42 PM
144
145Step 3: Compact
146The reconcile_view now contains the most up-to-date set of records and is now synchronized with changes from the RDBMS source system. For BI Reporting and Analytical tools, a reporting_table can be generated from the reconcile_view. Before creating this table, any previous instances of the table should be dropped as in the example below.
147
148reporting_table
149
150[code language=â€SQLâ€]
151DROP TABLE reporting_table;
152CREATE TABLE reporting_table AS
153SELECT * FROM reconcile_view;
154[/code]
155
156Moving the Reconciled View (reconcile_view) to a Reporting Table (reporting_table) reduces the amount of processing needed for reporting queries.
157
158Further, the data stored in the Reporting Table will also be static, unchanged until the next processing cycle. This provides consistency in reporting between processing cycles. In contrast, the Reconciled View (reconcile_view) is dynamic and will change as soon as new files (holding change records) are added to or removed from the Change table (incremental_table) folder /user/hive/incremental_table.
159
160Step 4: Purge
161To prepare for the next series of incremental records from the source, replace the Base table (base_table) with only the most up-to-date records (reporting_table). Also, delete the previously imported Change record content (incremental_table) by deleting the files located in the external table location (‘/user/hive/incremental_table’).
162
163From a HIVE client:
164[code language=â€SQLâ€]
165DROP TABLE base_table;
166CREATE TABLE base_table AS
167SELECT * FROM reporting_table;
168[/code]
169
170From an HDFS client:
171hadoop fs –rm –r /user/hive/incremental_table/*
172
173Final Thoughts
174While there are several possible approaches to supporting incremental data feeds into Hive, this example has a few key advantages:
175
176By maintaining an External Table for updates only, the table contents can be refreshed by simply adding or deleting files to that folder.
177The four steps in the processing cycle (Ingest, Reconcile, Compact and Purge) can be coordinated in a single OOZIE workflow. The OOZIE workflow can be a scheduled event that corresponds to the data freshness SLA (i.e. Daily, Weekly, Monthly, etc.)
178In addition to supporting INSERT and UPDATE synchronization, DELETES can be synchronized by adding either a DELETE_FLAG or DELETE_DATE field to the import source. Then, use this field as a filter in the Hive reporting table to hide deleted records. For example,
179[code language=â€SQLâ€]
180CREATE VIEW reconcile_view AS
181SELECT t1.* FROM
182(SELECT * FROM base_table
183UNION
184SELECT * FROM incremental_table) t1
185JOIN
186(SELECT id, max(modified_date) max_modified FROM
187(SELECT * FROM base_table
188UNION
189SELECT * FROM incremental_table)
190GROUP BY id) s
191ON t1.id = s.id AND t1.modified_date = s.max_modified
192AND t1.delete_date IS NULL;
193[/code]
194In essence, this four-step strategy enables incremental updates, as we await the near-term development
195of Hive support for full CRUD operations (Insert, Select, Update, Delete).