Reading multiple rows from sqlite and send to Erlang process in one bulk.
This commit is contained in:
@@ -54,7 +54,7 @@ typedef enum {
|
|||||||
cmd_changes,
|
cmd_changes,
|
||||||
cmd_prepare,
|
cmd_prepare,
|
||||||
cmd_bind,
|
cmd_bind,
|
||||||
cmd_step,
|
cmd_multi_step,
|
||||||
cmd_reset,
|
cmd_reset,
|
||||||
cmd_column_names,
|
cmd_column_names,
|
||||||
cmd_column_types,
|
cmd_column_types,
|
||||||
@@ -476,49 +476,66 @@ make_cell(ErlNifEnv *env, sqlite3_stmt *statement, unsigned int i)
|
|||||||
}
|
}
|
||||||
|
|
||||||
static ERL_NIF_TERM
|
static ERL_NIF_TERM
|
||||||
make_row(ErlNifEnv *env, sqlite3_stmt *statement)
|
make_row(ErlNifEnv *env, sqlite3_stmt *statement, ERL_NIF_TERM *array, int size)
|
||||||
{
|
{
|
||||||
int i, size;
|
|
||||||
ERL_NIF_TERM *array;
|
|
||||||
ERL_NIF_TERM row;
|
|
||||||
|
|
||||||
size = sqlite3_column_count(statement);
|
|
||||||
array = (ERL_NIF_TERM *) enif_alloc(sizeof(ERL_NIF_TERM)*size);
|
|
||||||
|
|
||||||
if(!array)
|
if(!array)
|
||||||
return make_error_tuple(env, "no_memory");
|
return make_error_tuple(env, "no_memory");
|
||||||
|
|
||||||
for(i = 0; i < size; i++)
|
for(int i = 0; i < size; i++)
|
||||||
array[i] = make_cell(env, statement, i);
|
array[i] = make_cell(env, statement, i);
|
||||||
|
|
||||||
row = make_row_tuple(env, enif_make_tuple_from_array(env, array, size));
|
return enif_make_tuple_from_array(env, array, size);
|
||||||
enif_free(array);
|
|
||||||
return row;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
static ERL_NIF_TERM
|
static ERL_NIF_TERM
|
||||||
do_step(ErlNifEnv *env, sqlite3 *db, sqlite3_stmt *stmt)
|
do_multi_step(ErlNifEnv *env, sqlite3 *db, sqlite3_stmt *stmt, const ERL_NIF_TERM arg)
|
||||||
{
|
{
|
||||||
|
ERL_NIF_TERM status;
|
||||||
|
ERL_NIF_TERM rows = enif_make_list_from_array(env, NULL, 0);
|
||||||
|
ERL_NIF_TERM *rowBuffer = NULL;
|
||||||
|
int rowBufferSize = 0;
|
||||||
|
|
||||||
|
int chunk_size = 0;
|
||||||
|
enif_get_int(env, arg, &chunk_size);
|
||||||
|
|
||||||
int rc = sqlite3_step(stmt);
|
int rc = sqlite3_step(stmt);
|
||||||
|
while (rc == SQLITE_ROW && chunk_size-- > 0)
|
||||||
|
{
|
||||||
|
if (!rowBufferSize)
|
||||||
|
rowBufferSize = sqlite3_column_count(stmt);
|
||||||
|
if (rowBuffer == NULL)
|
||||||
|
rowBuffer = (ERL_NIF_TERM *) enif_alloc(sizeof(ERL_NIF_TERM)*rowBufferSize);
|
||||||
|
|
||||||
if(rc == SQLITE_ROW)
|
rows = enif_make_list_cell(env, make_row(env, stmt, rowBuffer, rowBufferSize), rows);
|
||||||
return make_row(env, stmt);
|
|
||||||
if(rc == SQLITE_BUSY)
|
|
||||||
return make_atom(env, "$busy");
|
|
||||||
|
|
||||||
if(rc == SQLITE_DONE) {
|
if (chunk_size > 0)
|
||||||
/*
|
rc = sqlite3_step(stmt);
|
||||||
* Automatically reset the statement after a done so
|
|
||||||
* column_names will work after the statement is done.
|
|
||||||
*
|
|
||||||
* Not resetting the statement can lead to vm crashes.
|
|
||||||
*/
|
|
||||||
sqlite3_reset(stmt);
|
|
||||||
return make_atom(env, "$done");
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/* We use prepare_v2, so any error code can be returned. */
|
switch(rc) {
|
||||||
return make_sqlite3_error_tuple(env, rc, db);
|
case SQLITE_ROW:
|
||||||
|
status = make_atom(env, "rows");
|
||||||
|
break;
|
||||||
|
case SQLITE_BUSY:
|
||||||
|
status = make_atom(env, "$busy");
|
||||||
|
break;
|
||||||
|
case SQLITE_DONE:
|
||||||
|
/*
|
||||||
|
* Automatically reset the statement after a done so
|
||||||
|
* column_names will work after the statement is done.
|
||||||
|
*
|
||||||
|
* Not resetting the statement can lead to vm crashes.
|
||||||
|
*/
|
||||||
|
sqlite3_reset(stmt);
|
||||||
|
status = make_atom(env, "$done");
|
||||||
|
break;
|
||||||
|
default:
|
||||||
|
/* We use prepare_v2, so any error code can be returned. */
|
||||||
|
return make_sqlite3_error_tuple(env, rc, db);
|
||||||
|
}
|
||||||
|
|
||||||
|
enif_free(rowBuffer);
|
||||||
|
return enif_make_tuple2(env, status, rows);
|
||||||
}
|
}
|
||||||
|
|
||||||
static ERL_NIF_TERM
|
static ERL_NIF_TERM
|
||||||
@@ -630,8 +647,8 @@ evaluate_command(esqlite_command *cmd, esqlite_connection *conn)
|
|||||||
return do_changes(cmd->env, conn, cmd->arg);
|
return do_changes(cmd->env, conn, cmd->arg);
|
||||||
case cmd_prepare:
|
case cmd_prepare:
|
||||||
return do_prepare(cmd->env, conn, cmd->arg);
|
return do_prepare(cmd->env, conn, cmd->arg);
|
||||||
case cmd_step:
|
case cmd_multi_step:
|
||||||
return do_step(cmd->env, conn->db, stmt->statement);
|
return do_multi_step(cmd->env, conn->db, stmt->statement, cmd->arg);
|
||||||
case cmd_reset:
|
case cmd_reset:
|
||||||
return do_reset(cmd->env, conn->db, stmt->statement);
|
return do_reset(cmd->env, conn->db, stmt->statement);
|
||||||
case cmd_bind:
|
case cmd_bind:
|
||||||
@@ -916,39 +933,47 @@ esqlite_bind(ErlNifEnv *env, int argc, const ERL_NIF_TERM argv[])
|
|||||||
}
|
}
|
||||||
|
|
||||||
/*
|
/*
|
||||||
* Step to a prepared statement
|
* Multi step to a prepared statement
|
||||||
*/
|
*/
|
||||||
static ERL_NIF_TERM
|
static ERL_NIF_TERM
|
||||||
esqlite_step(ErlNifEnv *env, int argc, const ERL_NIF_TERM argv[])
|
esqlite_multi_step(ErlNifEnv *env, int argc, const ERL_NIF_TERM argv[])
|
||||||
{
|
{
|
||||||
esqlite_connection *conn;
|
esqlite_connection *conn;
|
||||||
esqlite_statement *stmt;
|
esqlite_statement *stmt;
|
||||||
esqlite_command *cmd = NULL;
|
esqlite_command *cmd = NULL;
|
||||||
ErlNifPid pid;
|
ErlNifPid pid;
|
||||||
|
int chunk_size = 0;
|
||||||
|
|
||||||
if(argc != 4)
|
if(argc != 5)
|
||||||
return enif_make_badarg(env);
|
return enif_make_badarg(env);
|
||||||
|
|
||||||
if(!enif_get_resource(env, argv[0], esqlite_connection_type, (void **) &conn))
|
if(!enif_get_resource(env, argv[0], esqlite_connection_type, (void **) &conn))
|
||||||
return enif_make_badarg(env);
|
return enif_make_badarg(env);
|
||||||
|
|
||||||
if(!enif_get_resource(env, argv[1], esqlite_statement_type, (void **) &stmt))
|
if(!enif_get_resource(env, argv[1], esqlite_statement_type, (void **) &stmt))
|
||||||
return enif_make_badarg(env);
|
return enif_make_badarg(env);
|
||||||
if(!enif_is_ref(env, argv[2]))
|
|
||||||
return make_error_tuple(env, "invalid_ref");
|
if(!enif_get_int(env, argv[2], &chunk_size))
|
||||||
if(!enif_get_local_pid(env, argv[3], &pid))
|
return make_error_tuple(env, "invalid_chunk_size");
|
||||||
return make_error_tuple(env, "invalid_pid");
|
|
||||||
|
if(!enif_is_ref(env, argv[3]))
|
||||||
|
return make_error_tuple(env, "invalid_ref");
|
||||||
|
|
||||||
|
if(!enif_get_local_pid(env, argv[4], &pid))
|
||||||
|
return make_error_tuple(env, "invalid_pid");
|
||||||
|
|
||||||
if(!stmt->statement)
|
if(!stmt->statement)
|
||||||
return make_error_tuple(env, "no_prepared_statement");
|
return make_error_tuple(env, "no_prepared_statement");
|
||||||
|
|
||||||
cmd = command_create();
|
cmd = command_create();
|
||||||
if(!cmd)
|
if(!cmd)
|
||||||
return make_error_tuple(env, "command_create_failed");
|
return make_error_tuple(env, "command_create_failed");
|
||||||
|
|
||||||
cmd->type = cmd_step;
|
cmd->type = cmd_multi_step;
|
||||||
cmd->ref = enif_make_copy(cmd->env, argv[2]);
|
cmd->ref = enif_make_copy(cmd->env, argv[3]);
|
||||||
cmd->pid = pid;
|
cmd->pid = pid;
|
||||||
cmd->stmt = enif_make_copy(cmd->env, argv[1]);
|
cmd->stmt = enif_make_copy(cmd->env, argv[1]);
|
||||||
|
cmd->arg = enif_make_copy(cmd->env, argv[2]);
|
||||||
|
|
||||||
return push_command(env, conn, cmd);
|
return push_command(env, conn, cmd);
|
||||||
}
|
}
|
||||||
@@ -1133,7 +1158,7 @@ static ErlNifFunc nif_funcs[] = {
|
|||||||
{"changes", 3, esqlite_changes},
|
{"changes", 3, esqlite_changes},
|
||||||
{"prepare", 4, esqlite_prepare},
|
{"prepare", 4, esqlite_prepare},
|
||||||
{"insert", 4, esqlite_insert},
|
{"insert", 4, esqlite_insert},
|
||||||
{"step", 4, esqlite_step},
|
{"multi_step", 5, esqlite_multi_step},
|
||||||
{"reset", 4, esqlite_reset},
|
{"reset", 4, esqlite_reset},
|
||||||
// TODO: {"esqlite_bind", 3, esqlite_bind_named},
|
// TODO: {"esqlite_bind", 3, esqlite_bind_named},
|
||||||
{"bind", 5, esqlite_bind},
|
{"bind", 5, esqlite_bind},
|
||||||
|
|||||||
120
src/esqlite3.erl
120
src/esqlite3.erl
@@ -31,6 +31,7 @@
|
|||||||
bind/2, bind/3,
|
bind/2, bind/3,
|
||||||
fetchone/1,
|
fetchone/1,
|
||||||
fetchall/1,
|
fetchall/1,
|
||||||
|
fetchall/2,
|
||||||
column_names/1, column_names/2,
|
column_names/1, column_names/2,
|
||||||
column_types/1, column_types/2,
|
column_types/1, column_types/2,
|
||||||
close/1, close/2]).
|
close/1, close/2]).
|
||||||
@@ -38,6 +39,7 @@
|
|||||||
-export([q/2, q/3, map/3, foreach/3]).
|
-export([q/2, q/3, map/3, foreach/3]).
|
||||||
|
|
||||||
-define(DEFAULT_TIMEOUT, 5000).
|
-define(DEFAULT_TIMEOUT, 5000).
|
||||||
|
-define(DEFAULT_CHUNK_SIZE, 5000).
|
||||||
|
|
||||||
%%
|
%%
|
||||||
-type connection() :: {connection, reference(), term()}.
|
-type connection() :: {connection, reference(), term()}.
|
||||||
@@ -123,19 +125,19 @@ foreach(F, Sql, Connection) ->
|
|||||||
Row :: tuple(),
|
Row :: tuple(),
|
||||||
ColumnNames :: tuple().
|
ColumnNames :: tuple().
|
||||||
foreach_s(F, Statement) when is_function(F, 1) ->
|
foreach_s(F, Statement) when is_function(F, 1) ->
|
||||||
case try_step(Statement, 0) of
|
case try_multi_step(Statement, 1, [], 0) of
|
||||||
'$done' -> ok;
|
{'$done', []} -> ok;
|
||||||
{error, _} = E -> F(E);
|
{error, _} = E -> F(E);
|
||||||
{row, Row} ->
|
{rows, [Row | []]} ->
|
||||||
F(Row),
|
F(Row),
|
||||||
foreach_s(F, Statement)
|
foreach_s(F, Statement)
|
||||||
end;
|
end;
|
||||||
foreach_s(F, Statement) when is_function(F, 2) ->
|
foreach_s(F, Statement) when is_function(F, 2) ->
|
||||||
ColumnNames = column_names(Statement),
|
ColumnNames = column_names(Statement),
|
||||||
case try_step(Statement, 0) of
|
case try_multi_step(Statement, 1, [], 0) of
|
||||||
'$done' -> ok;
|
{'$done', []} -> ok;
|
||||||
{error, _} = E -> F([], E);
|
{error, _} = E -> F([], E);
|
||||||
{row, Row} ->
|
{rows, [Row | []]} ->
|
||||||
F(ColumnNames, Row),
|
F(ColumnNames, Row),
|
||||||
foreach_s(F, Statement)
|
foreach_s(F, Statement)
|
||||||
end.
|
end.
|
||||||
@@ -147,59 +149,79 @@ foreach_s(F, Statement) when is_function(F, 2) ->
|
|||||||
ColumnNames :: tuple(),
|
ColumnNames :: tuple(),
|
||||||
Type :: term().
|
Type :: term().
|
||||||
map_s(F, Statement) when is_function(F, 1) ->
|
map_s(F, Statement) when is_function(F, 1) ->
|
||||||
case try_step(Statement, 0) of
|
case try_multi_step(Statement, 1, [], 0) of
|
||||||
'$done' -> [];
|
{'$done', []} -> [];
|
||||||
{error, _} = E -> F(E);
|
{error, _} = E -> F(E);
|
||||||
{row, Row} ->
|
{rows, [Row | []]} ->
|
||||||
[F(Row) | map_s(F, Statement)]
|
[F(Row) | map_s(F, Statement)]
|
||||||
end;
|
end;
|
||||||
map_s(F, Statement) when is_function(F, 2) ->
|
map_s(F, Statement) when is_function(F, 2) ->
|
||||||
ColumnNames = column_names(Statement),
|
ColumnNames = column_names(Statement),
|
||||||
case try_step(Statement, 0) of
|
case try_multi_step(Statement, 1, [], 0) of
|
||||||
'$done' -> [];
|
{'$done', []} -> [];
|
||||||
{error, _} = E -> F([], E);
|
{error, _} = E -> F([], E);
|
||||||
{row, Row} ->
|
{rows, [Row | []]} ->
|
||||||
[F(ColumnNames, Row) | map_s(F, Statement)]
|
[F(ColumnNames, Row) | map_s(F, Statement)]
|
||||||
end.
|
end.
|
||||||
|
|
||||||
%%
|
%%
|
||||||
-spec fetchone(statement()) -> tuple().
|
%%-spec fetchone(statement()) -> tuple().
|
||||||
fetchone(Statement) ->
|
fetchone(Statement) ->
|
||||||
case try_step(Statement, 0) of
|
case try_multi_step(Statement, 1, [], 0) of
|
||||||
'$done' -> ok;
|
{'$done', []} -> ok;
|
||||||
{error, _} = E -> E;
|
{error, _} = E -> E;
|
||||||
{row, Row} -> Row
|
{rows, [Row | []]} -> Row
|
||||||
end.
|
end.
|
||||||
|
|
||||||
%%
|
%% @doc Fetch all records
|
||||||
|
%% @param Statement is prepared sql statement
|
||||||
|
%% @spec fetchall(statement()) -> list(tuple()) | {error, term()}.
|
||||||
-spec fetchall(statement()) ->
|
-spec fetchall(statement()) ->
|
||||||
list(tuple()) |
|
list(tuple()) |
|
||||||
{error, term()}.
|
{error, term()}.
|
||||||
fetchall(Statement) ->
|
fetchall(Statement) ->
|
||||||
case try_step(Statement, 0) of
|
fetchall(Statement, ?DEFAULT_CHUNK_SIZE).
|
||||||
'$done' ->
|
|
||||||
[];
|
%% @doc Fetch all records
|
||||||
{error, _} = E -> E;
|
%% @param Statement is prepared sql statement
|
||||||
{row, Row} ->
|
%% @param ChunkSize is a count of rows to read from sqlite and send to erlang process in one bulk.
|
||||||
case fetchall(Statement) of
|
%% Decrease this value if rows are heavy. Default value is 5000 (DEFAULT_CHUNK_SIZE).
|
||||||
{error, _} = E -> E;
|
%% @spec fetchall(statement()) -> list(tuple()) | {error, term()}.
|
||||||
Rest -> [Row | Rest]
|
-spec fetchall(statement(), pos_integer()) ->
|
||||||
end
|
list(tuple()) |
|
||||||
|
{error, term()}.
|
||||||
|
fetchall(Statement, ChunkSize) ->
|
||||||
|
case fetchall_internal(Statement, ChunkSize, []) of
|
||||||
|
{'$done', Rows} -> lists:reverse(Rows);
|
||||||
|
{error, _} = E -> E
|
||||||
end.
|
end.
|
||||||
|
|
||||||
%% Try the step, when the database is busy,
|
%% return rows in revers order
|
||||||
-spec try_step(statement(), non_neg_integer()) ->
|
-spec fetchall_internal(statement(), pos_integer(), list(tuple())) ->
|
||||||
'$done' |
|
{'$done', list(tuple())} |
|
||||||
term().
|
{error, term()}.
|
||||||
try_step(_Statement, Tries) when Tries > 5 ->
|
fetchall_internal(Statement, ChunkSize, Rest) ->
|
||||||
|
case try_multi_step(Statement, ChunkSize, Rest, 0) of
|
||||||
|
{rows, Rows} -> fetchall_internal(Statement, ChunkSize, Rows);
|
||||||
|
Else -> Else
|
||||||
|
end.
|
||||||
|
|
||||||
|
%% Try a number of steps, when the database is busy,
|
||||||
|
%% return rows in revers order
|
||||||
|
-spec try_multi_step(statement(), pos_integer(), list(tuple()), non_neg_integer()) ->
|
||||||
|
{rows, list(tuple())} |
|
||||||
|
{'$done', list(tuple())} |
|
||||||
|
{error, term()}.
|
||||||
|
try_multi_step(_Statement, _ChunkSize, _Rest, Tries) when Tries > 5 ->
|
||||||
throw(too_many_tries);
|
throw(too_many_tries);
|
||||||
try_step(Statement, Tries) ->
|
try_multi_step(Statement, ChunkSize, Rest, Tries) ->
|
||||||
case esqlite3:step(Statement) of
|
case multi_step(Statement, ChunkSize) of
|
||||||
'$busy' ->
|
{'$busy', Rows} -> %% core can fetch a number of rows (rows < ChunkSize) per 'multi_step' call and then get busy...
|
||||||
|
erlang:display({"busy", Tries}),
|
||||||
timer:sleep(100 * Tries),
|
timer:sleep(100 * Tries),
|
||||||
try_step(Statement, Tries + 1);
|
try_multi_step(Statement, ChunkSize, Rows ++ Rest, Tries + 1);
|
||||||
Something ->
|
{Status, Rows} ->
|
||||||
Something
|
{Status, Rows ++ Rest}
|
||||||
end.
|
end.
|
||||||
|
|
||||||
%% @doc Execute Sql statement, returns the number of affected rows.
|
%% @doc Execute Sql statement, returns the number of affected rows.
|
||||||
@@ -280,7 +302,29 @@ step(Stmt) ->
|
|||||||
-spec step(term(), timeout()) -> tuple() | '$busy' | '$done'.
|
-spec step(term(), timeout()) -> tuple() | '$busy' | '$done'.
|
||||||
step({statement, Stmt, {connection, _, Conn}}, Timeout) ->
|
step({statement, Stmt, {connection, _, Conn}}, Timeout) ->
|
||||||
Ref = make_ref(),
|
Ref = make_ref(),
|
||||||
ok = esqlite3_nif:step(Conn, Stmt, Ref, self()),
|
ok = esqlite3_nif:multi_step(Conn, Stmt, 1, Ref, self()),
|
||||||
|
case receive_answer(Ref, Timeout) of
|
||||||
|
{rows, [Row | []]} -> {row, Row};
|
||||||
|
{'$done', []} -> '$done';
|
||||||
|
{'$busy', []} -> '$busy';
|
||||||
|
Else -> Else
|
||||||
|
end.
|
||||||
|
|
||||||
|
%% make multiple sqlite steps per call
|
||||||
|
%% return rows in reverse order
|
||||||
|
multi_step(Stmt, ChunkSize) ->
|
||||||
|
multi_step(Stmt, ChunkSize, ?DEFAULT_TIMEOUT).
|
||||||
|
|
||||||
|
%% make multiple sqlite steps per call
|
||||||
|
%% return rows in reverse order
|
||||||
|
-spec multi_step(term(), pos_integer(), timeout()) ->
|
||||||
|
{rows, list(tuple())} |
|
||||||
|
{'$busy', list(tuple())} |
|
||||||
|
{'$done', list(tuple())} |
|
||||||
|
{error, term()}.
|
||||||
|
multi_step({statement, Stmt, {connection, _, Conn}}, ChunkSize, Timeout) ->
|
||||||
|
Ref = make_ref(),
|
||||||
|
ok = esqlite3_nif:multi_step(Conn, Stmt, ChunkSize, Ref, self()),
|
||||||
receive_answer(Ref, Timeout).
|
receive_answer(Ref, Timeout).
|
||||||
|
|
||||||
%% @doc Reset the prepared statement back to its initial state.
|
%% @doc Reset the prepared statement back to its initial state.
|
||||||
|
|||||||
@@ -27,7 +27,7 @@
|
|||||||
changes/3,
|
changes/3,
|
||||||
insert/4,
|
insert/4,
|
||||||
prepare/4,
|
prepare/4,
|
||||||
step/4,
|
multi_step/5,
|
||||||
reset/4,
|
reset/4,
|
||||||
finalize/4,
|
finalize/4,
|
||||||
bind/5,
|
bind/5,
|
||||||
@@ -90,8 +90,8 @@ prepare(_Db, _Ref, _Dest, _Sql) ->
|
|||||||
|
|
||||||
%% @doc
|
%% @doc
|
||||||
%%
|
%%
|
||||||
%% @spec step(statement(), reference(), pid()) -> ok | {error, message()}
|
%% @spec multi_step(statement(), pos_integer(), reference(), pid()) -> {term(), list(tuple)} | {error, message()}
|
||||||
step(_Db, _Stmt, _Ref, _Dest) ->
|
multi_step(_Db, _Stmt, _Chunk_Size, _Ref, _Dest) ->
|
||||||
erlang:nif_error(nif_library_not_loaded).
|
erlang:nif_error(nif_library_not_loaded).
|
||||||
|
|
||||||
%% @doc
|
%% @doc
|
||||||
|
|||||||
Reference in New Issue
Block a user