如何使用Kiba-ETL将数据表转换为集合的哈希结构
Hey there! I get exactly what you're trying to do—turn each column (hash key) from your source data into an array of its unique values. Let's break down how to implement this cleanly in Kiba.
Core Approach
Since Kiba processes records row-by-row by default, we'll need to:
- First accumulate all incoming records in memory (note: this works best if your dataset isn't astronomically large; if it's huge, you'd want to use a database temp table or streaming aggregation, but let's focus on the common case first)
- Then, after collecting all records, iterate over each column (hash key) to extract and deduplicate its values
- Finally, output the aggregated result as a single record (or multiple records, depending on your needs)
Step-by-Step Implementation
1. Build a Record Accumulator Transformer
We'll create a simple transformer class that keeps track of all records as they flow through the pipeline:
class RecordAccumulator def initialize @records = [] end def process(row) @records << row # Return nil here so we don't pass rows further down the pipeline yet nil end attr_reader :records end
2. Add a Post-Processing Step to Aggregate Unique Values
Kiba lets you run code after all records are processed using post_process. We'll use this to take our accumulated records and build the unique value sets for each column:
accumulator = RecordAccumulator.new Kiba.parse do # Your source here (replace with your actual source, e.g., CSV, JSON, etc.) source MySource, "path/to/your/data" # Pass all rows through the accumulator transform accumulator # Post-process to build the unique value sets post_process do aggregated = {} # First, get all unique columns from the records all_columns = accumulator.records.flat_map(&:keys).uniq all_columns.each do |column| # Extract all values for the column, remove nils, deduplicate, convert to array aggregated[column] = accumulator.records.map { |row| row[column] }.compact.uniq end # Output the aggregated result (you could also write this to a destination here) puts aggregated.inspect # Or if you need to send it to a destination, you can yield it: # yield aggregated end # Optional: Add a destination if you want to write the aggregated result somewhere # destination MyDestination, "path/to/output" end
3. Example Output
For your sample input:
[ {dairy: "Milk", protein: "Steak", carb: "Potatoes"}, {dairy: "Milk", protein: "Eggs", carb: "Potatoes"}, {dairy: "Cheese", protein: "Steak", carb: "Potatoes"}, {dairy: "Cream", protein: "Eggs", carb: "Rice"} ]
The aggregated result would be:
{ dairy: ["Milk", "Cheese", "Cream"], protein: ["Steak", "Eggs"], carb: ["Potatoes", "Rice"] }
Notes for Larger Datasets
If your data is too big to fit in memory, you'll want to use a different approach:
- Use a database as an intermediate step: write all rows to a temp table, then run
SELECT DISTINCTqueries for each column - Use a streaming aggregation library (like Redis) to track unique values as rows are processed, instead of storing all rows in memory
That should get you exactly what you need! Let me know if you need help adapting this to your specific source/destination setup.
内容的提问来源于stack exchange,提问作者Gabriel Fortuna

