基于Pig的日期维度数据清洗:按地区与年份统计连接数需求
Alright, let's tackle this problem step by step—you've got two datasets to work with, need to compute cross-location connection counts per year, plus handle Pig data cleaning. Here's a practical, actionable approach:
First, we need to clean both datasets to remove invalid records and standardize formats, which is critical for accurate downstream calculations.
1.1 Clean ID-Location Dataset
The raw dataset has fields (ID, beginning_year, ending_year, location). Common issues include missing values, invalid year ranges, and empty location strings.
-- Load raw ID-location data (adjust delimiter if needed) location_raw = LOAD 'path/to/your/location_data' USING PigStorage(',') AS (id:chararray, begin_year:int, end_year:int, location:chararray); -- Filter out invalid records location_clean = FILTER location_raw BY id IS NOT NULL AND begin_year IS NOT NULL AND end_year IS NOT NULL AND location IS NOT NULL AND begin_year <= end_year; -- Ensure valid time range -- Deduplicate if same ID has multiple valid entries (keep unique records) location_dedup = DISTINCT location_clean; -- Store cleaned data for later use STORE location_dedup INTO 'path/to/cleaned_location_data' USING PigStorage(',');
1.2 Clean ID-Connection Dataset
The raw dataset has fields (ID1, ID2, connection_date). We'll convert dates to years (since we need yearly counts) and filter out invalid entries like self-connections.
-- Load raw connection data connection_raw = LOAD 'path/to/your/connection_data' USING PigStorage(',') AS (id1:chararray, id2:chararray, conn_date:chararray); -- Extract year from connection date (assuming date format is YYYY-MM-DD or YYYY) connection_with_year = FOREACH connection_raw GENERATE id1, id2, SUBSTRING(conn_date, 0, 4) AS conn_year:int; -- Filter invalid records connection_clean = FILTER connection_with_year BY id1 IS NOT NULL AND id2 IS NOT NULL AND conn_year IS NOT NULL AND id1 != id2; -- Exclude self-connections -- Deduplicate duplicate connection entries connection_dedup = DISTINCT connection_clean; -- Store cleaned data STORE connection_dedup INTO 'path/to/cleaned_connection_data' USING PigStorage(',');
Since connections never expire after creation, for each target year Y, we need to count all connections created on or before Y, where ID1 and ID2 were located in Loc1 and Loc2 during Y.
2.1 Generate Yearly Location Mapping
First, we'll expand each ID's location record to cover every year in its valid time range (from begin_year to end_year).
-- Load cleaned location data location_clean = LOAD 'path/to/cleaned_location_data' USING PigStorage(',') AS (id:chararray, begin_year:int, end_year:int, location:chararray); -- Generate all years between begin and end for each ID -- *Note: GENERATE_SEQUENCE requires Pig 0.11 or later* yearly_location = FOREACH location_clean GENERATE id, FLATTEN(GENERATE_SEQUENCE(begin_year, end_year)) AS year:int, location; -- Store yearly location mapping STORE yearly_location INTO 'path/to/yearly_location_mapping' USING PigStorage(',');
2.2 Join Connections with Yearly Locations
Now we'll link connections to the locations of ID1 and ID2 for every year after the connection was created.
-- Load cleaned connection data connection_clean = LOAD 'path/to/cleaned_connection_data' USING PigStorage(',') AS (id1:chararray, id2:chararray, conn_year:int); -- Load yearly location mapping yearly_location = LOAD 'path/to/yearly_location_mapping' USING PigStorage(',') AS (id:chararray, year:int, location:chararray); -- Join connections with ID1's location for all years >= connection creation year conn_id1_loc = JOIN connection_clean BY id1, yearly_location BY id; conn_id1_loc_filtered = FILTER conn_id1_loc BY yearly_location::year >= connection_clean::conn_year; -- Join with ID2's location for the same year final_join = JOIN conn_id1_loc_filtered BY (id2, yearly_location::year), yearly_location BY (id, year); -- Rename fields for clarity final_data = FOREACH final_join GENERATE conn_id1_loc_filtered::year AS year:int, conn_id1_loc_filtered::location AS location1:chararray, yearly_location::location AS location2:chararray, conn_id1_loc_filtered::id1 AS id1, conn_id1_loc_filtered::id2 AS id2;
2.3 Aggregate to Get Final Counts
Finally, group by year, location1, and location2 to count unique connections.
-- Group by year, location1, location2 and count distinct connections connection_counts = FOREACH (GROUP final_data BY (year, location1, location2)) GENERATE group.year AS year:int, group.location1 AS location1:chararray, group.location2 AS location2:chararray, COUNT(DISTINCT (final_data.id1, final_data.id2)) AS connection_count:int; -- Optional: Sort results for readability sorted_counts = ORDER connection_counts BY year ASC, location1 ASC, location2 ASC; -- Store the final result in your desired format STORE sorted_counts INTO 'path/to/final_connection_counts' USING PigStorage(',');
Key Notes
- If an ID has multiple location records for the same year (e.g., moved mid-year), you'll need to add logic to pick a primary location (e.g., earliest/latest entry).
- Adjust filters if you want to include/exclude same-location connections (
location1 == location2). - For older Pig versions without
GENERATE_SEQUENCE, use a custom UDF or precompute a year lookup table to expand the location records.
内容的提问来源于stack exchange,提问作者Claire Cui

