const std = @import("std");const os = std.os;const mem = std.mem;const assert = std.debug.assert;const stdout = std.io.getStdOut().writer();const Allocator = std.mem.Allocator;const ArrayList = std.ArrayList; var column_count: u64 = 103; // FIXMEvar row_count: u64 = 49369568; // FIXME fn extract_field(row: *[]u8, delim: u8) []u8 { var p: u64 = 0; while (p < row.len) { defer p += 1; if (row.*[p] == delim) { const result = row.*[0..p]; row.* = row.*[p + 1 ..]; return result; } } const result = row.*; row.* = row.*[row.len..]; return result;} const Column = struct { name: []u8, range: []u8, count: u64, metrics: Metrics, const Metrics = struct { unique_value_count: u64 = 0, min_length: u64 = std.math.maxInt(u64), max_length: u64 = std.math.minInt(u64), is_nullable: bool = false, is_unique: bool = true, is_constant: bool = true, is_integer: bool = true, is_uppercase_hex: bool = true, is_lowercase_hex: bool = true, }; const Iterator = struct { next_index: u64 = 0, count: u64, range: []u8, pub fn next(it: *Iterator) ?[]u8 { if (it.next_index == it.count) return null; it.next_index += 1; return extract_field(&it.range, '|'); } }; pub fn iterator(self: *const Column) Iterator { return .{ .range = self.range, .count = self.count }; }}; const AnalyzeColumnContext = struct { column: *Column, mutex: *std.Mutex,}; fn analyze_column(context: AnalyzeColumnContext) !void { var column = context.column; var mutex = context.mutex; var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); defer arena.deinit(); const allocator = &arena.allocator; var unique_values = std.StringHashMap(void).init(allocator); defer unique_values.deinit(); // TODO: write metrics to a temp object to reduce cache contention // (or pad/align `Column` to a multiple of 64 bytes?) column.metrics = .{}; var it = column.iterator(); while (it.next()) |field| { if (field.len == 0) { column.metrics.is_nullable = true; } else { try unique_values.put(field, {}); column.metrics.min_length = std.math.min(column.metrics.min_length, field.len); column.metrics.max_length = std.math.max(column.metrics.max_length, field.len); for (field) |ch| { const is_digit = ('0' <= ch and ch <= '9'); const is_lowercase_hex = ('0' <= ch and ch <= '9') or ('a' <= ch and ch <= 'f'); const is_uppercase_hex = ('0' <= ch and ch <= '9') or ('A' <= ch and ch <= 'F'); if (!is_digit) { column.metrics.is_integer = false; } if (!is_lowercase_hex) { column.metrics.is_lowercase_hex = false; } if (!is_uppercase_hex) { column.metrics.is_uppercase_hex = false; } } } } // TODO: Separate values that appear only once/more than once? // That would reduce memory consumption during data implication analysis. column.metrics.unique_value_count = unique_values.count(); column.metrics.is_unique = (column.metrics.unique_value_count == row_count); column.metrics.is_constant = (column.metrics.unique_value_count == 1); { const lock = mutex.acquire(); defer lock.release(); // FIXME: probably do this somewhere else (then we don't need the mutex) try stdout.print("{}:\n", .{column.name}); if (column.metrics.is_integer) { try stdout.print(" integer\n", .{}); } else if (column.metrics.is_uppercase_hex) { try stdout.print(" uppercase hex\n", .{}); } else if (column.metrics.is_lowercase_hex) { try stdout.print(" lowercase hex\n", .{}); } else { try stdout.print(" string\n", .{}); } if (column.metrics.is_unique) { try stdout.print(" unique\n", .{}); } if (column.metrics.is_constant) { try stdout.print(" constant\n", .{}); } if (column.metrics.is_nullable) { try stdout.print(" nullable\n", .{}); } try stdout.print(" unique values: {}\n", .{column.metrics.unique_value_count}); if (column.metrics.min_length == column.metrics.max_length) { try stdout.print(" length: {}\n", .{column.metrics.min_length}); } try stdout.print("\n", .{}); }} ////////////////////////////////////////////////////////////////////////////// // TODO: convert to tagged union?const DataCondition = struct { column_name: []const u8, value: []const u8, is_equality: bool, pub fn print(self: *const DataCondition) !void { try stdout.print("{}", .{self.column_name}); if (self.is_equality) { try stdout.print(" == ", .{}); } else { try stdout.print(" != ", .{}); } if (self.value.len == 0) { try stdout.print("null", .{}); } else { try stdout.print("\"{}\"", .{self.value}); } } pub fn eql(self: DataCondition, other: DataCondition) bool { if (!mem.eql(u8, self.column_name, other.column_name)) return false; if (!mem.eql(u8, self.value, other.value)) return false; if (self.is_equality != other.is_equality) return false; return true; }}; const DataImplication = struct { antecedent: DataCondition, consequent: DataCondition, is_equivalence: bool, pub fn print(self: *const DataImplication) !void { try self.antecedent.print(); if (self.is_equivalence) { try stdout.print(" <=> ", .{}); } else { try stdout.print(" ==> ", .{}); } try self.consequent.print(); }}; const AnalyzeDataImplicationContext = struct { column_1: *const Column, column_2: *const Column, null_values: std.StringHashMap(void), non_null_values: std.StringHashMap(void), null_non_values: std.StringHashMap(void), non_null_non_values: std.StringHashMap(void), pub fn init(allocator: *Allocator, column_1: *const Column, column_2: *const Column) AnalyzeDataImplicationContext { return AnalyzeDataImplicationContext{ .column_1 = column_1, .column_2 = column_2, .null_values = std.StringHashMap(void).init(allocator), .non_null_values = std.StringHashMap(void).init(allocator), .null_non_values = std.StringHashMap(void).init(allocator), .non_null_non_values = std.StringHashMap(void).init(allocator), }; } pub fn deinit(self: *AnalyzeDataImplicationContext) void { self.null_values.deinit(); self.non_null_values.deinit(); self.null_non_values.deinit(); self.non_null_non_values.deinit(); } pub fn needs_update(self: *AnalyzeDataImplicationContext) bool { if (!self.column_1.metrics.is_nullable) return false; if (self.null_values.count() < 2) return true; if (self.non_null_non_values.count() < 2) return true; if (self.non_null_values.count() < 2) return true; if (self.null_non_values.count() < 2) return true; return false; } pub fn update(self: *AnalyzeDataImplicationContext, field_1: []const u8, field_2: []const u8) !void { if (!self.column_1.metrics.is_nullable) return; if (field_1.len == 0) { if (self.null_values.count() < 2) { try self.null_values.put(field_2, {}); } if (self.non_null_non_values.count() < 2) { try self.non_null_non_values.put(field_2, {}); } } else { if (self.non_null_values.count() < 2) { try self.non_null_values.put(field_2, {}); } if (self.null_non_values.count() < 2) { try self.null_non_values.put(field_2, {}); } } } pub fn finalize(self: *AnalyzeDataImplicationContext, output: []DataImplication) u64 { if (!self.column_1.metrics.is_nullable) return 0; var output_count: u64 = 0; if (self.null_values.count() == 1) { const field_2 = blk: { var it = self.null_values.iterator(); if (it.next()) |entry| { break :blk entry.key; } else unreachable; }; output[output_count] = DataImplication{ .antecedent = DataCondition{ .column_name = self.column_2.name, .value = field_2, .is_equality = false, }, .consequent = DataCondition{ .column_name = self.column_1.name, .value = "", .is_equality = false, }, .is_equivalence = false, }; output_count += 1; } if (self.null_non_values.count() == 1) { const field_2 = blk: { var it = self.null_non_values.iterator(); if (it.next()) |entry| { break :blk entry.key; } else unreachable; }; output[output_count] = DataImplication{ .antecedent = DataCondition{ .column_name = self.column_2.name, .value = field_2, .is_equality = true, }, .consequent = DataCondition{ .column_name = self.column_1.name, .value = "", .is_equality = false, }, .is_equivalence = false, }; output_count += 1; } if (self.non_null_values.count() == 1) { const field_2 = blk: { var it = self.non_null_values.iterator(); if (it.next()) |entry| { break :blk entry.key; } else unreachable; }; output[output_count] = DataImplication{ .antecedent = DataCondition{ .column_name = self.column_2.name, .value = field_2, .is_equality = false, }, .consequent = DataCondition{ .column_name = self.column_1.name, .value = "", .is_equality = true, }, .is_equivalence = false, }; output_count += 1; } if (self.non_null_non_values.count() == 1) { const field_2 = blk: { var it = self.non_null_non_values.iterator(); if (it.next()) |entry| { break :blk entry.key; } else unreachable; }; output[output_count] = DataImplication{ .antecedent = DataCondition{ .column_name = self.column_2.name, .value = field_2, .is_equality = true, }, .consequent = DataCondition{ .column_name = self.column_1.name, .value = "", .is_equality = true, }, .is_equivalence = false, }; output_count += 1; } return output_count; }}; const AnalyzeColumnPairContext = struct { column_1: *Column, column_2: *Column, mutex: *std.Mutex,}; fn analyze_column_pair(context: AnalyzeColumnPairContext) !void { var column_1 = context.column_1; var column_2 = context.column_2; var mutex = context.mutex; var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); defer arena.deinit(); const allocator = &arena.allocator; // Don't report trivial implications. if (column_1.metrics.is_unique or column_1.metrics.is_constant) { return; } if (column_2.metrics.is_unique or column_2.metrics.is_constant) { return; } var impl_1 = true; var impl_2 = true; if (column_1.metrics.unique_value_count < column_2.metrics.unique_value_count) { impl_1 = false; } if (column_2.metrics.unique_value_count < column_1.metrics.unique_value_count) { impl_2 = false; } var function_1 = std.StringHashMap([]u8).init(allocator); defer function_1.deinit(); if (impl_1) { var capacity = column_1.metrics.unique_value_count; // TODO: use HashMap method capacity += capacity >> 2; capacity += 1; try function_1.ensureCapacity(@intCast(u32, capacity)); } var function_2 = std.StringHashMap([]u8).init(allocator); defer function_2.deinit(); if (impl_2) { var capacity = column_2.metrics.unique_value_count; // TODO: use HashMap method capacity += capacity >> 2; capacity += 1; try function_2.ensureCapacity(@intCast(u32, capacity)); } var analyze_data_implication_context_1 = AnalyzeDataImplicationContext.init(allocator, column_1, column_2); defer analyze_data_implication_context_1.deinit(); var analyze_data_implication_context_2 = AnalyzeDataImplicationContext.init(allocator, column_2, column_1); defer analyze_data_implication_context_2.deinit(); var it_1 = column_1.iterator(); var it_2 = column_2.iterator(); { var i: u64 = 0; while (i < row_count and (impl_1 or impl_2 or analyze_data_implication_context_1.needs_update() or analyze_data_implication_context_2.needs_update())) { defer i += 1; const field_1 = blk: { if (it_1.next()) |field_1| { break :blk field_1; } else unreachable; }; const field_2 = blk: { if (it_2.next()) |field_2| { break :blk field_2; } else unreachable; }; if (impl_1) { const maybe_old_entry = try function_1.fetchPut(field_1, field_2); if (maybe_old_entry) |old_entry| { if (!mem.eql(u8, old_entry.value, field_2)) { impl_1 = false; } } } if (impl_2) { const maybe_old_entry = try function_2.fetchPut(field_2, field_1); if (maybe_old_entry) |old_entry| { if (!mem.eql(u8, old_entry.value, field_1)) { impl_2 = false; } } } try analyze_data_implication_context_1.update(field_1, field_2); try analyze_data_implication_context_2.update(field_2, field_1); } } // FIXME: probably do this somewhere else (then we don't need the mutex) { const lock = mutex.acquire(); defer lock.release(); var backing_buffer_1: [8]DataImplication = undefined; var backing_buffer_2: [8]DataImplication = undefined; var buffer: *[8]DataImplication = &backing_buffer_1; var new_buffer: *[8]DataImplication = &backing_buffer_2; var data_implications: []DataImplication = undefined; { var p: u64 = 0; p += analyze_data_implication_context_1.finalize(buffer[p..]); p += analyze_data_implication_context_2.finalize(buffer[p..]); data_implications = buffer[0..p]; } { // Simplify "A ==> B", "B ==> A" to "A <=> B". var was_processed: [8]bool = undefined; mem.set(bool, was_processed[0..], false); var p: u64 = 0; blk: for (data_implications) |data_implication, i| { if (was_processed[i]) continue; defer was_processed[i] = true; for (data_implications) |other_data_implication, j| { if (was_processed[j]) continue; if (DataCondition.eql(data_implication.antecedent, other_data_implication.consequent) and DataCondition.eql(data_implication.consequent, other_data_implication.antecedent) and data_implication.is_equivalence == other_data_implication.is_equivalence) { defer was_processed[j] = true; new_buffer[p] = DataImplication{ .antecedent = data_implication.antecedent, .consequent = data_implication.consequent, .is_equivalence = true, }; p += 1; continue :blk; } } new_buffer[p] = data_implication; p += 1; } data_implications = new_buffer[0..p]; mem.swap(*[8]DataImplication, &buffer, &new_buffer); } { // Simplify "A ==> B", "~A ==> ~B" to "A <=> B". var was_processed: [8]bool = undefined; mem.set(bool, was_processed[0..], false); var p: u64 = 0; blk: for (data_implications) |data_implication, i| { if (was_processed[i]) continue; defer was_processed[i] = true; if (data_implication.is_equivalence == false) { for (data_implications) |other_data_implication, j| { if (was_processed[j]) continue; if (other_data_implication.is_equivalence == false) { if (mem.eql(u8, data_implication.antecedent.column_name, other_data_implication.antecedent.column_name) and mem.eql(u8, data_implication.antecedent.value, other_data_implication.antecedent.value) and data_implication.antecedent.is_equality != other_data_implication.antecedent.is_equality and mem.eql(u8, data_implication.consequent.column_name, other_data_implication.consequent.column_name) and mem.eql(u8, data_implication.consequent.value, other_data_implication.consequent.value) and data_implication.consequent.is_equality != other_data_implication.consequent.is_equality) { defer was_processed[j] = true; new_buffer[p] = DataImplication{ .antecedent = data_implication.antecedent, .consequent = data_implication.consequent, .is_equivalence = true, }; p += 1; continue :blk; } } } } new_buffer[p] = data_implication; p += 1; } data_implications = new_buffer[0..p]; mem.swap(*[8]DataImplication, &buffer, &new_buffer); } { // Simplify "A <=> B", "~A <=> ~B" to "A <=> B". var was_processed: [8]bool = undefined; mem.set(bool, was_processed[0..], false); var p: u64 = 0; blk: for (data_implications) |data_implication, i| { if (was_processed[i]) continue; defer was_processed[i] = true; if (data_implication.is_equivalence == true) { for (data_implications) |other_data_implication, j| { if (was_processed[j]) continue; if (other_data_implication.is_equivalence == true) { if (mem.eql(u8, data_implication.antecedent.column_name, other_data_implication.antecedent.column_name) and mem.eql(u8, data_implication.antecedent.value, other_data_implication.antecedent.value) and data_implication.antecedent.is_equality != other_data_implication.antecedent.is_equality and mem.eql(u8, data_implication.consequent.column_name, other_data_implication.consequent.column_name) and mem.eql(u8, data_implication.consequent.value, other_data_implication.consequent.value) and data_implication.consequent.is_equality != other_data_implication.consequent.is_equality) { defer was_processed[j] = true; new_buffer[p] = data_implication; p += 1; continue :blk; } } // TODO: also handle "A <=> B", "~B <=> ~A" case. } } new_buffer[p] = data_implication; p += 1; } data_implications = new_buffer[0..p]; mem.swap(*[8]DataImplication, &buffer, &new_buffer); } // TODO: canonicalize remaining implications for (data_implications) |data_implication, i| { try data_implication.print(); try stdout.print("\n", .{}); } if (impl_1 and impl_2) { try stdout.print("{} <=> {}\n", .{ column_1.name, column_2.name }); } else if (impl_1) { try stdout.print("{} ==> {}\n", .{ column_1.name, column_2.name }); } else if (impl_2) { try stdout.print("{} ==> {}\n", .{ column_2.name, column_1.name }); } else {} }} pub fn main() !void { var gpa = std.heap.GeneralPurposeAllocator(.{}){}; defer assert(!gpa.deinit()); const allocator = &gpa.allocator; const fd = try os.open("matricula.transposed.csv", os.O_RDONLY, 0); defer os.close(fd); try os.lseek_END(fd, 0); const length = try os.lseek_CUR_get(fd); const data = try os.mmap( null, length, os.PROT_READ, os.MAP_SHARED, fd, 0, ); defer os.munmap(data); var columns = try allocator.alloc(Column, column_count); defer allocator.free(columns); var input_stream = data; const header = extract_field(&input_stream, '\n'); var header_stream = header; for (columns) |*column| { column.name = extract_field(&header_stream, '|'); column.range = extract_field(&input_stream, '\n'); column.count = row_count; // FIXME } { var mutex = std.Mutex{}; for (columns) |*column, j| { var thread = try std.Thread.spawn(AnalyzeColumnContext{ .column = column, .mutex = &mutex }, analyze_column); // TODO: set up thread pool etc. thread.wait(); } } if (true) { var mutex = std.Mutex{}; for (columns) |*column_1, j| { for (columns[(j + 1)..]) |*column_2| { var thread = try std.Thread.spawn(AnalyzeColumnPairContext{ .column_1 = column_1, .column_2 = column_2, .mutex = &mutex }, analyze_column_pair); // TODO: set up thread pool etc. thread.wait(); } } }}
Comments