Skip to content

Commit 9251676

Browse files
committed
Add Arrow CSR query APIs
1 parent e318e85 commit 9251676

13 files changed

Lines changed: 541 additions & 4 deletions

package-lock.json

Lines changed: 193 additions & 3 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

package.json

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,7 @@
3939
"tmp": "^0.2.3"
4040
},
4141
"dependencies": {
42+
"apache-arrow": "^21.1.0",
4243
"cmake-js": "^8.0.0",
4344
"node-addon-api": "^6.0.0"
4445
}

src_cpp/include/node_connection.h

Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -32,8 +32,10 @@ class NodeConnection : public Napi::ObjectWrap<NodeConnection> {
3232
void SetQueryTimeout(const Napi::CallbackInfo& info);
3333
Napi::Value ExecuteAsync(const Napi::CallbackInfo& info);
3434
Napi::Value QueryAsync(const Napi::CallbackInfo& info);
35+
Napi::Value QueryArrowAsync(const Napi::CallbackInfo& info);
3536
Napi::Value ExecuteSync(const Napi::CallbackInfo& info);
3637
Napi::Value QuerySync(const Napi::CallbackInfo& info);
38+
Napi::Value QueryArrowSync(const Napi::CallbackInfo& info);
3739
Napi::Value CreateArrowTableSync(const Napi::CallbackInfo& info);
3840
Napi::Value CreateArrowRelTableSync(const Napi::CallbackInfo& info);
3941
Napi::Value DropArrowTableSync(const Napi::CallbackInfo& info);
@@ -188,5 +190,40 @@ class ConnectionQueryAsyncWorker : public Napi::AsyncWorker {
188190
std::optional<Napi::ThreadSafeFunction> progressCallback;
189191
};
190192

193+
class ConnectionQueryArrowAsyncWorker : public Napi::AsyncWorker {
194+
public:
195+
ConnectionQueryArrowAsyncWorker(Napi::Function& callback,
196+
std::shared_ptr<Connection>& connection, std::shared_ptr<Database>& database,
197+
std::string statement, int64_t chunkSize, NodeQueryResult* nodeQueryResult)
198+
: Napi::AsyncWorker(callback), connection(connection), database(database),
199+
statement(std::move(statement)), chunkSize(chunkSize), nodeQueryResult(nodeQueryResult) {}
200+
201+
~ConnectionQueryArrowAsyncWorker() override = default;
202+
203+
void Execute() override {
204+
try {
205+
auto result = connection->queryAsArrow(statement, chunkSize);
206+
if (!result->isSuccess()) {
207+
SetError(result->getErrorMessage());
208+
return;
209+
}
210+
nodeQueryResult->AdoptQueryResult(std::move(result), connection, database);
211+
} catch (const std::exception& exc) {
212+
SetError(std::string(exc.what()));
213+
}
214+
}
215+
216+
void OnOK() override { Callback().Call({Env().Null()}); }
217+
218+
void OnError(Napi::Error const& error) override { Callback().Call({error.Value()}); }
219+
220+
private:
221+
std::shared_ptr<Connection> connection;
222+
std::shared_ptr<Database> database;
223+
std::string statement;
224+
int64_t chunkSize;
225+
NodeQueryResult* nodeQueryResult;
226+
};
227+
191228
} // namespace main
192229
} // namespace lbug

src_cpp/include/node_query_result.h

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,7 @@ class NodeQueryResult : public Napi::ObjectWrap<NodeQueryResult> {
4343
Napi::Value GetColumnNamesSync(const Napi::CallbackInfo& info);
4444
Napi::Value GetQuerySummarySync(const Napi::CallbackInfo& info);
4545
Napi::Value GetQuerySummaryAsync(const Napi::CallbackInfo& info);
46+
Napi::Value GetCSRSync(const Napi::CallbackInfo& info);
4647
void PopulateColumnNames();
4748
void Close(const Napi::CallbackInfo& info);
4849
void Close();

0 commit comments

Comments
 (0)