MCPcopy Create free account
hub / github.com/BIT-DataLab/LakeBench / sub_process

Function sub_process

union/D3L/query.py:34–73  ·  view source on GitHub ↗
(query_tables, queue, output_folder, qe, dataloader, k_value)

Source from the content-addressed store, hash-verified

32 return result
33
34def sub_process(query_tables, queue, output_folder, qe, dataloader, k_value):
35 for i, table_name_with_extension in enumerate(query_tables):
36
37 table_name = os.path.splitext(table_name_with_extension)[0]
38 if not os.path.exists(output_folder):
39 os.makedirs(output_folder)
40
41 results_file = os.path.join(output_folder, f"{table_name}.csv")
42
43 if os.path.exists(results_file):
44 print(f"跳过查询表 {i + 1},因为结果文件已经存在:{results_file}")
45 queue.put(1) # 在队列中放入一个占位符值
46 continue
47
48 # 执行查询
49 results, extended_results = qe.table_query(table=dataloader.read_table(table_name=table_name),
50 aggregator=None, k=k_value, verbose=True)
51
52 # 创建一个新的 CSV 文件
53 with open(results_file, mode='w', newline='') as csvfile:
54 writer = csv.writer(csvfile)
55
56 # Write the header
57 writer.writerow(['query_table', 'candidate_table', 'query_col_name', 'candidate_col_name'])
58
59 # 写入查询结果到 CSV 文件
60 for result in extended_results:
61 query_table = table_name
62 candidate_table = os.path.basename(result[0])
63
64 for x, column_info in enumerate(result[1]):
65 column_info_name = f"column_info{x+1}"
66 globals()[column_info_name] = column_info[0]
67
68 row = [f"{query_table}.csv", f"{candidate_table}.csv", column_info[0][0], column_info[0][1]]
69 writer.writerow(row)
70
71 print(f"Results for query table {i + 1} have been written to {results_file}")
72 queue.put(1)
73 queue.put((-1, "test-pid"))
74
75def merge_csv(output_folder, combined_file_path):
76 csv_files = glob.glob(output_folder + "/*.csv")

Callers

nothing calls this directly

Calls 2

table_queryMethod · 0.80
read_tableMethod · 0.45

Tested by

no test coverage detected