Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
Metz
oceanbase
提交
93fd1a9e
O
oceanbase
项目概览
Metz
/
oceanbase
与 Fork 源项目一致
Fork自
oceanbase / oceanbase
通知
1
Star
0
Fork
0
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
0
列表
看板
标记
里程碑
合并请求
0
DevOps
流水线
流水线任务
计划
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
O
oceanbase
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
0
Issue
0
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
DevOps
DevOps
流水线
流水线任务
计划
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
流水线任务
提交
Issue看板
提交
93fd1a9e
编写于
3月 14, 2022
作者:
X
xy0
提交者:
LINGuanRen
3月 14, 2022
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
[asan] 3.1 open observer无法启动
上级
ac5edad6
变更
5
隐藏空白更改
内联
并排
Showing
5 changed file
with
140 addition
and
126 deletion
+140
-126
src/observer/mysql/ob_sync_cmd_driver.cpp
src/observer/mysql/ob_sync_cmd_driver.cpp
+17
-16
src/observer/mysql/ob_sync_plan_driver.cpp
src/observer/mysql/ob_sync_plan_driver.cpp
+23
-22
src/sql/engine/prepare/ob_execute_executor.cpp
src/sql/engine/prepare/ob_execute_executor.cpp
+72
-70
src/sql/parser/sql_parser_mysql_mode.y
src/sql/parser/sql_parser_mysql_mode.y
+26
-16
src/storage/ob_sstable_multi_version_row_iterator.cpp
src/storage/ob_sstable_multi_version_row_iterator.cpp
+2
-2
未找到文件。
src/observer/mysql/ob_sync_cmd_driver.cpp
浏览文件 @
93fd1a9e
...
...
@@ -28,15 +28,15 @@ using namespace obmysql;
using
namespace
share
;
namespace
observer
{
ObSyncCmdDriver
::
ObSyncCmdDriver
(
const
ObGlobalContext
&
gctx
,
const
ObSqlCtx
&
ctx
,
sql
::
ObSQLSessionInfo
&
session
,
ObQueryRetryCtrl
&
retry_ctrl
,
ObIMPPacketSender
&
sender
)
ObSyncCmdDriver
::
ObSyncCmdDriver
(
const
ObGlobalContext
&
gctx
,
const
ObSqlCtx
&
ctx
,
sql
::
ObSQLSessionInfo
&
session
,
ObQueryRetryCtrl
&
retry_ctrl
,
ObIMPPacketSender
&
sender
)
:
ObQueryDriver
(
gctx
,
ctx
,
session
,
retry_ctrl
,
sender
)
{}
ObSyncCmdDriver
::~
ObSyncCmdDriver
()
{}
int
ObSyncCmdDriver
::
response_result
(
ObMySQLResultSet
&
result
)
int
ObSyncCmdDriver
::
response_result
(
ObMySQLResultSet
&
result
)
{
int
ret
=
OB_SUCCESS
;
bool
process_ok
=
false
;
...
...
@@ -77,7 +77,7 @@ int ObSyncCmdDriver::response_result(ObMySQLResultSet& result)
}
}
else
{
OMPKEOF
eofp
;
const
ObWarningBuffer
*
warnings_buf
=
common
::
ob_get_tsi_warning_buffer
();
const
ObWarningBuffer
*
warnings_buf
=
common
::
ob_get_tsi_warning_buffer
();
uint16_t
warning_count
=
0
;
if
(
OB_ISNULL
(
warnings_buf
))
{
LOG_WARN
(
"can not get thread warnings buffer"
);
...
...
@@ -119,10 +119,10 @@ int ObSyncCmdDriver::response_result(ObMySQLResultSet& result)
}
else
if
(
!
result
.
is_with_rows
()
||
(
sender_
.
need_send_extra_ok_packet
()
&&
!
result
.
has_more_result
()))
{
process_ok
=
true
;
ObOKPParam
ok_param
;
ok_param
.
message_
=
const_cast
<
char
*>
(
result
.
get_message
());
ok_param
.
message_
=
const_cast
<
char
*>
(
result
.
get_message
());
ok_param
.
affected_rows_
=
result
.
get_affected_rows
();
ok_param
.
lii_
=
result
.
get_last_insert_id_to_client
();
const
ObWarningBuffer
*
warnings_buf
=
common
::
ob_get_tsi_warning_buffer
();
const
ObWarningBuffer
*
warnings_buf
=
common
::
ob_get_tsi_warning_buffer
();
if
(
OB_ISNULL
(
warnings_buf
))
{
LOG_WARN
(
"can not get thread warnings buffer"
);
}
else
{
...
...
@@ -152,7 +152,7 @@ int ObSyncCmdDriver::response_result(ObMySQLResultSet& result)
// two aspects:
// - set session last_schema_version to proxy for part DDL
// - promote local schema up to target version if last_schema_version is set
int
ObSyncCmdDriver
::
process_schema_version_changes
(
const
ObMySQLResultSet
&
result
)
int
ObSyncCmdDriver
::
process_schema_version_changes
(
const
ObMySQLResultSet
&
result
)
{
int
ret
=
OB_SUCCESS
;
...
...
@@ -171,11 +171,11 @@ int ObSyncCmdDriver::process_schema_version_changes(const ObMySQLResultSet& resu
if
(
OB_SUCC
(
ret
))
{
// - promote local schema up to target version if last_schema_version is set
if
(
result
.
get_stmt_type
()
==
stmt
::
T_VARIABLE_SET
)
{
const
ObVariableSetStmt
*
set_stmt
=
static_cast
<
const
ObVariableSetStmt
*>
(
result
.
get_cmd
());
const
ObVariableSetStmt
*
set_stmt
=
static_cast
<
const
ObVariableSetStmt
*>
(
result
.
get_cmd
());
if
(
NULL
!=
set_stmt
)
{
ObVariableSetStmt
::
VariableNamesSetNode
tmp_node
;
// just for init node
for
(
int64_t
i
=
0
;
OB_SUCC
(
ret
)
&&
i
<
set_stmt
->
get_variables_size
();
++
i
)
{
ObVariableSetStmt
::
VariableNamesSetNode
&
var_node
=
tmp_node
;
ObVariableSetStmt
::
VariableNamesSetNode
&
var_node
=
tmp_node
;
ObString
set_var_name
(
OB_SV_LAST_SCHEMA_VERSION
);
if
(
OB_FAIL
(
set_stmt
->
get_variable_node
(
i
,
var_node
)))
{
LOG_WARN
(
"fail to get_variable_node"
,
K
(
i
),
K
(
ret
));
...
...
@@ -218,12 +218,12 @@ int ObSyncCmdDriver::check_and_refresh_schema(uint64_t tenant_id)
return
ret
;
}
int
ObSyncCmdDriver
::
response_query_result
(
ObMySQLResultSet
&
result
)
int
ObSyncCmdDriver
::
response_query_result
(
ObMySQLResultSet
&
result
)
{
int
ret
=
OB_SUCCESS
;
session_
.
get_trans_desc
().
consistency_wait
();
const
common
::
ObNewRow
*
row
=
NULL
;
const
common
::
ObNewRow
*
row
=
NULL
;
if
(
OB_FAIL
(
result
.
next_row
(
row
)))
{
LOG_WARN
(
"fail to get next row"
,
K
(
ret
));
}
else
if
(
OB_FAIL
(
response_query_header
(
result
,
result
.
has_more_result
(),
true
)))
{
...
...
@@ -232,9 +232,9 @@ int ObSyncCmdDriver::response_query_result(ObMySQLResultSet& result)
ret
=
OB_ERR_UNEXPECTED
;
LOG_WARN
(
"session info is null"
,
K
(
ret
));
}
else
{
ObNewRow
*
tmp_row
=
const_cast
<
ObNewRow
*>
(
row
);
ObNewRow
*
tmp_row
=
const_cast
<
ObNewRow
*>
(
row
);
for
(
int64_t
i
=
0
;
OB_SUCC
(
ret
)
&&
i
<
tmp_row
->
get_count
();
i
++
)
{
ObObj
&
value
=
tmp_row
->
get_cell
(
i
);
ObObj
&
value
=
tmp_row
->
get_cell
(
i
);
if
(
ob_is_string_type
(
value
.
get_type
())
&&
CS_TYPE_INVALID
!=
value
.
get_collation_type
())
{
OZ
(
convert_string_value_charset
(
value
,
result
));
}
else
if
(
value
.
is_clob_locator
()
&&
OB_FAIL
(
convert_lob_value_charset
(
value
,
result
)))
{
...
...
@@ -247,14 +247,15 @@ int ObSyncCmdDriver::response_query_result(ObMySQLResultSet& result)
if
(
OB_SUCC
(
ret
))
{
MYSQL_PROTOCOL_TYPE
protocol_type
=
result
.
is_ps_protocol
()
?
BINARY
:
TEXT
;
const
ObSQLSessionInfo
*
tmp_session
=
result
.
get_exec_context
().
get_my_session
();
const
ObSQLSessionInfo
*
tmp_session
=
result
.
get_exec_context
().
get_my_session
();
const
ObDataTypeCastParams
dtc_params
=
ObBasicSessionInfo
::
create_dtc_params
(
tmp_session
);
O
MPKRow
rp
(
ObSMR
ow
(
protocol_type
,
O
bSMRow
sm_r
ow
(
protocol_type
,
*
row
,
dtc_params
,
result
.
get_field_columns
(),
ctx_
.
schema_guard_
,
tmp_session
->
get_effective_tenant_id
()));
tmp_session
->
get_effective_tenant_id
());
OMPKRow
rp
(
sm_row
);
if
(
OB_FAIL
(
sender_
.
response_packet
(
rp
)))
{
LOG_WARN
(
"response packet fail"
,
K
(
ret
),
KP
(
row
));
}
...
...
src/observer/mysql/ob_sync_plan_driver.cpp
浏览文件 @
93fd1a9e
...
...
@@ -30,22 +30,22 @@ using namespace sql;
using
namespace
obmysql
;
namespace
observer
{
ObSyncPlanDriver
::
ObSyncPlanDriver
(
const
ObGlobalContext
&
gctx
,
const
ObSqlCtx
&
ctx
,
sql
::
ObSQLSessionInfo
&
session
,
ObQueryRetryCtrl
&
retry_ctrl
,
ObIMPPacketSender
&
sender
)
ObSyncPlanDriver
::
ObSyncPlanDriver
(
const
ObGlobalContext
&
gctx
,
const
ObSqlCtx
&
ctx
,
sql
::
ObSQLSessionInfo
&
session
,
ObQueryRetryCtrl
&
retry_ctrl
,
ObIMPPacketSender
&
sender
)
:
ObQueryDriver
(
gctx
,
ctx
,
session
,
retry_ctrl
,
sender
)
{}
ObSyncPlanDriver
::~
ObSyncPlanDriver
()
{}
int
ObSyncPlanDriver
::
response_result
(
ObMySQLResultSet
&
result
)
int
ObSyncPlanDriver
::
response_result
(
ObMySQLResultSet
&
result
)
{
int
ret
=
OB_SUCCESS
;
bool
process_ok
=
false
;
// for select SQL
bool
ac
=
true
;
bool
admission_fail_and_need_retry
=
false
;
const
ObNewRow
*
not_used_row
=
NULL
;
const
ObNewRow
*
not_used_row
=
NULL
;
if
(
OB_ISNULL
(
result
.
get_physical_plan
()))
{
ret
=
OB_NOT_INIT
;
LOG_WARN
(
"should have set plan to result set"
,
K
(
ret
));
...
...
@@ -98,7 +98,7 @@ int ObSyncPlanDriver::response_result(ObMySQLResultSet& result)
process_ok
=
true
;
OMPKEOF
eofp
;
const
ObWarningBuffer
*
warnings_buf
=
common
::
ob_get_tsi_warning_buffer
();
const
ObWarningBuffer
*
warnings_buf
=
common
::
ob_get_tsi_warning_buffer
();
uint16_t
warning_count
=
0
;
if
(
OB_ISNULL
(
warnings_buf
))
{
LOG_WARN
(
"can not get thread warnings buffer"
);
...
...
@@ -149,10 +149,10 @@ int ObSyncPlanDriver::response_result(ObMySQLResultSet& result)
if
(
!
result
.
has_implicit_cursor
())
{
// no implicit cursor, send one ok packet to client
ObOKPParam
ok_param
;
ok_param
.
message_
=
const_cast
<
char
*>
(
result
.
get_message
());
ok_param
.
message_
=
const_cast
<
char
*>
(
result
.
get_message
());
ok_param
.
affected_rows_
=
result
.
get_affected_rows
();
ok_param
.
lii_
=
result
.
get_last_insert_id_to_client
();
const
ObWarningBuffer
*
warnings_buf
=
common
::
ob_get_tsi_warning_buffer
();
const
ObWarningBuffer
*
warnings_buf
=
common
::
ob_get_tsi_warning_buffer
();
if
(
OB_ISNULL
(
warnings_buf
))
{
LOG_WARN
(
"can not get thread warnings buffer"
);
}
else
{
...
...
@@ -169,7 +169,7 @@ int ObSyncPlanDriver::response_result(ObMySQLResultSet& result)
result
.
reset_implicit_cursor_idx
();
while
(
OB_SUCC
(
ret
)
&&
OB_SUCC
(
result
.
switch_implicit_cursor
()))
{
ObOKPParam
ok_param
;
ok_param
.
message_
=
const_cast
<
char
*>
(
result
.
get_message
());
ok_param
.
message_
=
const_cast
<
char
*>
(
result
.
get_message
());
ok_param
.
affected_rows_
=
result
.
get_affected_rows
();
ok_param
.
is_partition_hit_
=
session_
.
partition_hit
().
get_bool
();
ok_param
.
has_more_result_
=
!
result
.
is_cursor_end
();
...
...
@@ -201,12 +201,12 @@ int ObSyncPlanDriver::response_result(ObMySQLResultSet& result)
}
int
ObSyncPlanDriver
::
response_query_result
(
ObResultSet
&
result
,
bool
has_more_result
,
bool
&
can_retry
,
int64_t
fetch_limit
)
ObResultSet
&
result
,
bool
has_more_result
,
bool
&
can_retry
,
int64_t
fetch_limit
)
{
int
ret
=
OB_SUCCESS
;
can_retry
=
true
;
bool
is_first_row
=
true
;
const
ObNewRow
*
result_row
=
NULL
;
const
ObNewRow
*
result_row
=
NULL
;
bool
has_top_limit
=
result
.
get_has_top_limit
();
bool
is_cac_found_rows
=
result
.
is_calc_found_rows
();
int64_t
limit_count
=
OB_INVALID_COUNT
==
fetch_limit
?
INT64_MAX
:
fetch_limit
;
...
...
@@ -219,7 +219,7 @@ int ObSyncPlanDriver::response_query_result(
}
session_
.
get_trans_desc
().
consistency_wait
();
MYSQL_PROTOCOL_TYPE
protocol_type
=
result
.
is_ps_protocol
()
?
BINARY
:
TEXT
;
const
common
::
ColumnsFieldIArray
*
fields
=
NULL
;
const
common
::
ColumnsFieldIArray
*
fields
=
NULL
;
if
(
OB_SUCC
(
ret
))
{
fields
=
result
.
get_field_columns
();
if
(
OB_ISNULL
(
fields
))
{
...
...
@@ -228,7 +228,7 @@ int ObSyncPlanDriver::response_query_result(
}
}
while
(
OB_SUCC
(
ret
)
&&
row_num
<
limit_count
&&
!
OB_FAIL
(
result
.
get_next_row
(
result_row
)))
{
ObNewRow
*
row
=
const_cast
<
ObNewRow
*>
(
result_row
);
ObNewRow
*
row
=
const_cast
<
ObNewRow
*>
(
result_row
);
// If it is the first line, first reply to the client with field and other information
if
(
is_first_row
)
{
is_first_row
=
false
;
...
...
@@ -238,7 +238,7 @@ int ObSyncPlanDriver::response_query_result(
}
}
for
(
int64_t
i
=
0
;
OB_SUCC
(
ret
)
&&
i
<
row
->
get_count
();
i
++
)
{
ObObj
&
value
=
row
->
get_cell
(
i
);
ObObj
&
value
=
row
->
get_cell
(
i
);
if
(
result
.
is_ps_protocol
())
{
if
(
value
.
get_type
()
!=
fields
->
at
(
i
).
type_
.
get_type
())
{
ObCastCtx
cast_ctx
(
&
result
.
get_mem_pool
(),
NULL
,
CM_WARN_ON_FAIL
,
CS_TYPE_INVALID
);
...
...
@@ -259,12 +259,13 @@ int ObSyncPlanDriver::response_query_result(
}
if
(
OB_SUCC
(
ret
))
{
const
ObDataTypeCastParams
dtc_params
=
ObBasicSessionInfo
::
create_dtc_params
(
&
session_
);
O
MPKRow
rp
(
ObSMR
ow
(
protocol_type
,
O
bSMRow
sm_r
ow
(
protocol_type
,
*
row
,
dtc_params
,
result
.
get_field_columns
(),
ctx_
.
schema_guard_
,
session_
.
get_effective_tenant_id
()));
session_
.
get_effective_tenant_id
());
OMPKRow
rp
(
sm_row
);
if
(
OB_FAIL
(
sender_
.
response_packet
(
rp
)))
{
LOG_WARN
(
"response packet fail"
,
K
(
ret
),
KP
(
row
),
K
(
row_num
),
K
(
can_retry
));
// break;
...
...
@@ -297,12 +298,12 @@ int ObSyncPlanDriver::response_query_result(
return
ret
;
}
ObRemotePlanDriver
::
ObRemotePlanDriver
(
const
ObGlobalContext
&
gctx
,
const
ObSqlCtx
&
ctx
,
sql
::
ObSQLSessionInfo
&
session
,
ObQueryRetryCtrl
&
retry_ctrl
,
ObIMPPacketSender
&
sender
)
ObRemotePlanDriver
::
ObRemotePlanDriver
(
const
ObGlobalContext
&
gctx
,
const
ObSqlCtx
&
ctx
,
sql
::
ObSQLSessionInfo
&
session
,
ObQueryRetryCtrl
&
retry_ctrl
,
ObIMPPacketSender
&
sender
)
:
ObSyncPlanDriver
(
gctx
,
ctx
,
session
,
retry_ctrl
,
sender
)
{}
int
ObRemotePlanDriver
::
response_result
(
ObMySQLResultSet
&
result
)
int
ObRemotePlanDriver
::
response_result
(
ObMySQLResultSet
&
result
)
{
int
ret
=
result
.
get_errcode
();
bool
process_ok
=
false
;
...
...
@@ -347,7 +348,7 @@ int ObRemotePlanDriver::response_result(ObMySQLResultSet& result)
if
(
OB_SUCC
(
ret
))
{
process_ok
=
true
;
OMPKEOF
eofp
;
const
ObWarningBuffer
*
warnings_buf
=
common
::
ob_get_tsi_warning_buffer
();
const
ObWarningBuffer
*
warnings_buf
=
common
::
ob_get_tsi_warning_buffer
();
uint16_t
warning_count
=
0
;
if
(
OB_ISNULL
(
warnings_buf
))
{
LOG_WARN
(
"can not get thread warnings buffer"
);
...
...
@@ -392,10 +393,10 @@ int ObRemotePlanDriver::response_result(ObMySQLResultSet& result)
}
else
if
(
!
result
.
has_implicit_cursor
())
{
// no implicit cursor, send one ok packet to client
ObOKPParam
ok_param
;
ok_param
.
message_
=
const_cast
<
char
*>
(
result
.
get_message
());
ok_param
.
message_
=
const_cast
<
char
*>
(
result
.
get_message
());
ok_param
.
affected_rows_
=
result
.
get_affected_rows
();
ok_param
.
lii_
=
result
.
get_last_insert_id_to_client
();
const
ObWarningBuffer
*
warnings_buf
=
common
::
ob_get_tsi_warning_buffer
();
const
ObWarningBuffer
*
warnings_buf
=
common
::
ob_get_tsi_warning_buffer
();
if
(
OB_ISNULL
(
warnings_buf
))
{
LOG_WARN
(
"can not get thread warnings buffer"
);
}
else
{
...
...
@@ -412,7 +413,7 @@ int ObRemotePlanDriver::response_result(ObMySQLResultSet& result)
result
.
reset_implicit_cursor_idx
();
while
(
OB_SUCC
(
ret
)
&&
OB_SUCC
(
result
.
switch_implicit_cursor
()))
{
ObOKPParam
ok_param
;
ok_param
.
message_
=
const_cast
<
char
*>
(
result
.
get_message
());
ok_param
.
message_
=
const_cast
<
char
*>
(
result
.
get_message
());
ok_param
.
affected_rows_
=
result
.
get_affected_rows
();
ok_param
.
is_partition_hit_
=
session_
.
partition_hit
().
get_bool
();
ok_param
.
has_more_result_
=
!
result
.
is_cursor_end
();
...
...
src/sql/engine/prepare/ob_execute_executor.cpp
浏览文件 @
93fd1a9e
...
...
@@ -23,7 +23,7 @@ namespace oceanbase {
using
namespace
common
;
namespace
sql
{
int
ObExecuteExecutor
::
execute
(
ObExecContext
&
ctx
,
ObExecuteStmt
&
stmt
)
int
ObExecuteExecutor
::
execute
(
ObExecContext
&
ctx
,
ObExecuteStmt
&
stmt
)
{
int
ret
=
OB_SUCCESS
;
if
(
OB_ISNULL
(
ctx
.
get_sql_ctx
())
||
OB_ISNULL
(
ctx
.
get_my_session
()))
{
...
...
@@ -51,81 +51,83 @@ int ObExecuteExecutor::execute(ObExecContext& ctx, ObExecuteStmt& stmt)
}
if
(
OB_SUCC
(
ret
))
{
SMART_VAR
(
ObResultSet
,
result_set
,
*
ctx
.
get_my_session
())
observer
::
ObVirtualTableIteratorFactory
vt_iter_factory
(
*
GCTX
.
vt_iter_creator_
);
SMART_VAR
(
ObSqlCtx
,
sql_ctx
)
{
result_set
.
init_partition_location_cache
(
GCTX
.
location_cache_
,
GCTX
.
self_addr_
,
ctx
.
get_sql_ctx
()
->
schema_guard_
);
result_set
.
set_ps_protocol
();
ObTaskExecutorCtx
*
task_ctx
=
result_set
.
get_exec_context
().
get_task_executor_ctx
();
int64_t
tenant_version
=
0
;
int64_t
sys_version
=
0
;
if
(
OB_ISNULL
(
task_ctx
))
{
ret
=
OB_ERR_UNEXPECTED
;
LOG_ERROR
(
"task executor ctx can not be NULL"
,
K
(
task_ctx
),
K
(
ret
));
}
else
if
(
OB_FAIL
(
GCTX
.
schema_service_
->
get_tenant_received_broadcast_version
(
ctx
.
get_my_session
()
->
get_effective_tenant_id
(),
tenant_version
)))
{
LOG_WARN
(
"fail get tenant schema version"
,
K
(
ret
));
}
else
if
(
OB_FAIL
(
GCTX
.
schema_service_
->
get_tenant_received_broadcast_version
(
OB_SYS_TENANT_ID
,
sys_version
)))
{
LOG_WARN
(
"fail get sys schema version"
,
K
(
ret
));
}
else
{
ObSqlCtx
sql_ctx
;
observer
::
ObVirtualTableIteratorFactory
vt_iter_factory
(
*
GCTX
.
vt_iter_creator_
);
sql_ctx
.
retry_times_
=
0
;
sql_ctx
.
merged_version_
=
ctx
.
get_sql_ctx
()
->
merged_version_
;
sql_ctx
.
vt_iter_factory_
=
&
vt_iter_factory
;
sql_ctx
.
session_info_
=
ctx
.
get_my_session
();
sql_ctx
.
sql_proxy_
=
ctx
.
get_sql_ctx
()
->
sql_proxy_
;
sql_ctx
.
use_plan_cache_
=
true
;
sql_ctx
.
disable_privilege_check_
=
true
;
sql_ctx
.
partition_table_operator_
=
ctx
.
get_sql_ctx
()
->
partition_table_operator_
;
sql_ctx
.
partition_location_cache_
=
&
(
result_set
.
get_partition_location_cache
());
sql_ctx
.
part_mgr_
=
ctx
.
get_sql_ctx
()
->
part_mgr_
;
sql_ctx
.
is_prepare_protocol_
=
ctx
.
get_sql_ctx
()
->
is_prepare_protocol_
;
sql_ctx
.
is_prepare_stage_
=
ctx
.
get_sql_ctx
()
->
is_prepare_stage_
;
sql_ctx
.
schema_guard_
=
ctx
.
get_sql_ctx
()
->
schema_guard_
;
task_ctx
->
schema_service_
=
GCTX
.
schema_service_
;
task_ctx
->
set_query_tenant_begin_schema_version
(
tenant_version
);
task_ctx
->
set_query_sys_begin_schema_version
(
sys_version
);
task_ctx
->
set_min_cluster_version
(
GET_MIN_CLUSTER_VERSION
());
if
(
OB_FAIL
(
result_set
.
init
()))
{
LOG_WARN
(
"result set init failed"
,
K
(
ret
));
}
else
if
(
OB_FAIL
(
GCTX
.
sql_engine_
->
stmt_execute
(
stmt
.
get_prepare_id
(),
stmt
.
get_prepare_type
(),
params_array
,
sql_ctx
,
result_set
,
false
/* is_inner_sql */
)))
{
LOG_WARN
(
"failed to prepare stmt"
,
K
(
stmt
.
get_prepare_id
()),
K
(
stmt
.
get_prepare_type
()),
K
(
ret
));
SMART_VAR
(
ObResultSet
,
result_set
,
*
ctx
.
get_my_session
())
{
result_set
.
init_partition_location_cache
(
GCTX
.
location_cache_
,
GCTX
.
self_addr_
,
ctx
.
get_sql_ctx
()
->
schema_guard_
);
result_set
.
set_ps_protocol
();
ObTaskExecutorCtx
*
task_ctx
=
result_set
.
get_exec_context
().
get_task_executor_ctx
();
int64_t
tenant_version
=
0
;
int64_t
sys_version
=
0
;
if
(
OB_ISNULL
(
task_ctx
))
{
ret
=
OB_ERR_UNEXPECTED
;
LOG_ERROR
(
"task executor ctx can not be NULL"
,
K
(
task_ctx
),
K
(
ret
));
}
else
if
(
OB_FAIL
(
GCTX
.
schema_service_
->
get_tenant_received_broadcast_version
(
ctx
.
get_my_session
()
->
get_effective_tenant_id
(),
tenant_version
)))
{
LOG_WARN
(
"fail get tenant schema version"
,
K
(
ret
));
}
else
if
(
OB_FAIL
(
GCTX
.
schema_service_
->
get_tenant_received_broadcast_version
(
OB_SYS_TENANT_ID
,
sys_version
)))
{
LOG_WARN
(
"fail get sys schema version"
,
K
(
ret
));
}
else
{
if
(
OB_ISNULL
(
ctx
.
get_sql_ctx
()
->
schema_guard_
))
{
ret
=
OB_ERR_UNEXPECTED
;
LOG_WARN
(
"schema guard is null"
);
}
else
if
(
OB_FAIL
(
ctx
.
get_my_session
()
->
update_query_sensitive_system_variable
(
*
(
ctx
.
get_sql_ctx
()
->
schema_guard_
))))
{
LOG_WARN
(
"update query affacted system variable failed"
,
K
(
ret
));
}
else
if
(
OB_FAIL
(
result_set
.
sync_open
()))
{
LOG_WARN
(
"result set open failed"
,
K
(
result_set
.
get_statement_id
()),
K
(
ret
));
}
if
(
OB_SUCC
(
ret
))
{
if
(
result_set
.
is_with_rows
())
{
while
(
OB_SUCC
(
ret
))
{
const
common
::
ObNewRow
*
row
=
NULL
;
if
(
OB_FAIL
(
result_set
.
get_next_row
(
row
)))
{
if
(
OB_ITER_END
==
ret
)
{
ret
=
OB_SUCCESS
;
}
else
{
LOG_WARN
(
"get next row error"
,
K
(
ret
));
sql_ctx
.
retry_times_
=
0
;
sql_ctx
.
merged_version_
=
ctx
.
get_sql_ctx
()
->
merged_version_
;
sql_ctx
.
vt_iter_factory_
=
&
vt_iter_factory
;
sql_ctx
.
session_info_
=
ctx
.
get_my_session
();
sql_ctx
.
sql_proxy_
=
ctx
.
get_sql_ctx
()
->
sql_proxy_
;
sql_ctx
.
use_plan_cache_
=
true
;
sql_ctx
.
disable_privilege_check_
=
true
;
sql_ctx
.
partition_table_operator_
=
ctx
.
get_sql_ctx
()
->
partition_table_operator_
;
sql_ctx
.
partition_location_cache_
=
&
(
result_set
.
get_partition_location_cache
());
sql_ctx
.
part_mgr_
=
ctx
.
get_sql_ctx
()
->
part_mgr_
;
sql_ctx
.
is_prepare_protocol_
=
ctx
.
get_sql_ctx
()
->
is_prepare_protocol_
;
sql_ctx
.
is_prepare_stage_
=
ctx
.
get_sql_ctx
()
->
is_prepare_stage_
;
sql_ctx
.
schema_guard_
=
ctx
.
get_sql_ctx
()
->
schema_guard_
;
task_ctx
->
schema_service_
=
GCTX
.
schema_service_
;
task_ctx
->
set_query_tenant_begin_schema_version
(
tenant_version
);
task_ctx
->
set_query_sys_begin_schema_version
(
sys_version
);
task_ctx
->
set_min_cluster_version
(
GET_MIN_CLUSTER_VERSION
());
if
(
OB_FAIL
(
result_set
.
init
()))
{
LOG_WARN
(
"result set init failed"
,
K
(
ret
));
}
else
if
(
OB_FAIL
(
GCTX
.
sql_engine_
->
stmt_execute
(
stmt
.
get_prepare_id
(),
stmt
.
get_prepare_type
(),
params_array
,
sql_ctx
,
result_set
,
false
/* is_inner_sql */
)))
{
LOG_WARN
(
"failed to prepare stmt"
,
K
(
stmt
.
get_prepare_id
()),
K
(
stmt
.
get_prepare_type
()),
K
(
ret
));
}
else
{
if
(
OB_ISNULL
(
ctx
.
get_sql_ctx
()
->
schema_guard_
))
{
ret
=
OB_ERR_UNEXPECTED
;
LOG_WARN
(
"schema guard is null"
);
}
else
if
(
OB_FAIL
(
ctx
.
get_my_session
()
->
update_query_sensitive_system_variable
(
*
(
ctx
.
get_sql_ctx
()
->
schema_guard_
))))
{
LOG_WARN
(
"update query affacted system variable failed"
,
K
(
ret
));
}
else
if
(
OB_FAIL
(
result_set
.
sync_open
()))
{
LOG_WARN
(
"result set open failed"
,
K
(
result_set
.
get_statement_id
()),
K
(
ret
));
}
if
(
OB_SUCC
(
ret
))
{
if
(
result_set
.
is_with_rows
())
{
while
(
OB_SUCC
(
ret
))
{
const
common
::
ObNewRow
*
row
=
NULL
;
if
(
OB_FAIL
(
result_set
.
get_next_row
(
row
)))
{
if
(
OB_ITER_END
==
ret
)
{
ret
=
OB_SUCCESS
;
}
else
{
LOG_WARN
(
"get next row error"
,
K
(
ret
));
}
break
;
}
break
;
}
}
}
}
int
tmp_ret
=
OB_SUCCESS
;
if
((
tmp_ret
=
result_set
.
close
())
!=
OB_SUCCESS
)
{
LOG_WARN
(
"result set open failed"
,
K
(
result_set
.
get_statement_id
()),
K
(
ret
))
;
ret
=
OB_SUCCESS
==
ret
?
tmp_ret
:
ret
;
int
tmp_ret
=
OB_SUCCESS
;
if
((
tmp_ret
=
result_set
.
close
())
!=
OB_SUCCESS
)
{
LOG_WARN
(
"result set open failed"
,
K
(
result_set
.
get_statement_id
()),
K
(
ret
));
ret
=
OB_SUCCESS
==
ret
?
tmp_ret
:
ret
;
}
}
}
}
...
...
src/sql/parser/sql_parser_mysql_mode.y
浏览文件 @
93fd1a9e
...
...
@@ -9004,14 +9004,19 @@ table_reference inner_join_type opt_full_table_factor %prec LOWER_ON
if ($1->type_ == T_ORG) {
ParseNode *name_node = NULL;
make_name_node(name_node, result->malloc_pool_, "full");
malloc_non_terminal_node($$, result->malloc_pool_, T_ALIAS, $1->num_child_ + 1);
for (int i = 0; i <= $1->num_child_; ++i) {
if (i == 0) {
$$->children_[i] = $1->children_[i];
} else if (i == 1) {
$$->children_[i] = name_node;
} else {
$$->children_[i] = $1->children_[i - 1];
$$ = new_node(result->malloc_pool_, T_ALIAS, $1->num_child_ + 1);
if (OB_UNLIKELY($$ == NULL)) {
yyerror(NULL, result, "No more space for malloc\n");
YYABORT_NO_MEMORY;
} else {
for (int i = 0; i <= $1->num_child_; ++i) {
if (i == 0) {
$$->children_[i] = $1->children_[i];
} else if (i == 1) {
$$->children_[i] = name_node;
} else {
$$->children_[i] = $1->children_[i - 1];
}
}
}
} else if ($1->type_ == T_ALIAS && $1->children_[1] != NULL &&
...
...
@@ -9046,14 +9051,19 @@ table_factor %prec LOWER_COMMA
if ($1->type_ == T_ORG) {
ParseNode *name_node = NULL;
make_name_node(name_node, result->malloc_pool_, "full");
malloc_non_terminal_node($$, result->malloc_pool_, T_ALIAS, $1->num_child_ + 1);
for (int i = 0; i <= $1->num_child_; ++i) {
if (i == 0) {
$$->children_[i] = $1->children_[i];
} else if (i == 1) {
$$->children_[i] = name_node;
} else {
$$->children_[i] = $1->children_[i - 1];
$$ = new_node(result->malloc_pool_, T_ALIAS, $1->num_child_ + 1);
if (OB_UNLIKELY($$ == NULL)) {
yyerror(NULL, result, "No more space for malloc\n");
YYABORT_NO_MEMORY;
} else {
for (int i = 0; i <= $1->num_child_; ++i) {
if (i == 0) {
$$->children_[i] = $1->children_[i];
} else if (i == 1) {
$$->children_[i] = name_node;
} else {
$$->children_[i] = $1->children_[i - 1];
}
}
}
} else if ($1->type_ == T_ALIAS && $1->children_[1] != NULL &&
...
...
src/storage/ob_sstable_multi_version_row_iterator.cpp
浏览文件 @
93fd1a9e
...
...
@@ -306,7 +306,7 @@ int ObSSTableMultiVersionRowMultiGetter::inner_open(
}
}
if
(
OB_FAIL
(
ret
))
{
}
else
if
(
OB_FAIL
(
new_iterator
<
ObSSTableRowMultiScanner
>
(
*
access_ctx
.
allocator_
)))
{
}
else
if
(
OB_FAIL
(
new_iterator
<
ObSSTableRowMultiScanner
>
(
*
access_ctx
.
stmt_
allocator_
)))
{
LOG_WARN
(
"failed to new iterator"
,
K
(
ret
));
}
else
if
(
OB_FAIL
(
iter_
->
init
(
iter_param
,
access_ctx
,
table
,
&
multi_version_ranges_
)))
{
LOG_WARN
(
"failed to open multi scanner"
,
K
(
ret
));
...
...
@@ -431,7 +431,7 @@ int ObSSTableMultiVersionRowMultiScanner::inner_open(
}
if
(
OB_FAIL
(
ret
))
{
}
else
if
(
OB_FAIL
(
new_iterator
<
ObSSTableRowMultiScanner
>
(
*
access_ctx
.
allocator_
)))
{
}
else
if
(
OB_FAIL
(
new_iterator
<
ObSSTableRowMultiScanner
>
(
*
access_ctx
.
stmt_
allocator_
)))
{
LOG_WARN
(
"failed to new iterator"
,
K
(
ret
));
}
else
if
(
OB_FAIL
(
iter_
->
init
(
iter_param
,
access_ctx
,
table
,
&
multi_version_ranges_
)))
{
LOG_WARN
(
"failed to open scanner"
,
K
(
ret
));
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录