DTStack / DTStack/chunjun

Logminer+Oracle. Sink接收到数据并处理后,数据未插入目标库 并没有报错

Open
#679 0 comments 0 reactions 0 assignees View on GitHub
bug
Dominant language
Java
Stars
4.1k
Forks
1.7k
PR merge metrics
No merged PRs in 30d

Description

**Describe the bug**
A clear and concise description of what the bug is.

在Source 表插入数据后 ,sink端接收到数据 并处理完成,没有报错
但是 sink 表没有数据

![1649247106042](https://user-images.githubusercontent.com/49392769/161971739-a30e75cc-0a7a-40b5-8942-bf02d2158db1.jpg)

Chunjun SQL script:
+++++++++++++++++++++++++++++++++++++++++++++++
CREATE TABLE source_INP_PATIENT
(
HOS_CODE STRING,
PID STRING,
PKID STRING,
PATIENT_ID STRING,
HC_CARDNO STRING,
VISIT_CARDNO STRING,
INPNO STRING,
CASEID STRING,
NAME STRING,
BIRTHDAY DATE,
AGE DECIMAL,
AGE_DESC STRING,
SEX_NAME STRING,
IDCARD_NO STRING,
ADMISS_NUM DECIMAL,
INDAYS DECIMAL,
ADMISSION_DATE DATE,
DISCHARGE_DATE DATE,
STATUS DECIMAL,
ADMISS_WARD_CODE STRING,
ADMISS_DEPT_CODE STRING,
CURR_WARD_CODE STRING,
CURR_DEPT_CODE STRING,
CURR_BED_CODE STRING,
DISCHARGE_DEPT_CODE STRING,
DISCHARGE_WARD_CODE STRING,
CHIEF_DOCTOR_CODE STRING,
ATTENDING_DOCTOR_CODE STRING,
INP_DOCTOR_CODE STRING,
TEAM_CODE STRING,
INSURANCE_NATURE_CODE STRING,
INSURANCE_NATURE_NAME STRING,
INSURANCE_PLACE_CODE STRING,
INSURANCE_PLACE_NAME STRING,
INSURANCE_TYPE_CODE STRING,
INSURANCE_TYPE_NAME STRING,
MEDICAL_DIVISION_NAME STRING,
IS_DRG DECIMAL,
MODIFY_DATE DATE,
IS_DELETE DECIMAL,
IS_ADD_POINT DECIMAL,
DRG_TYPE DECIMAL,
ODS_CREATE_TIME DATE,
ODS_UPDATE_TIME DATE
)
WITH (
'connector' = 'oraclelogminer-x'
,'url' = 'jdbc:oracle:thin:@//xxxx:1521/helowin'
,'username' = 'xxx'
,'password' = 'xxx'
,'cat' = 'INSERT,UPDATE,DELETE'
,'table' = 'FLINKUSER.INP_PATIENT'
,'read-position' = 'current'
);

CREATE TABLE sink_INP_PATIENT_22
(
HOS_CODE STRING,
PID STRING,
PKID STRING,
PATIENT_ID STRING,
HC_CARDNO STRING,
VISIT_CARDNO STRING,
INPNO STRING,
CASEID STRING,
NAME STRING,
BIRTHDAY DATE,
AGE DECIMAL,
AGE_DESC STRING,
SEX_NAME STRING,
IDCARD_NO STRING,
ADMISS_NUM DECIMAL,
INDAYS DECIMAL,
ADMISSION_DATE DATE,
DISCHARGE_DATE DATE,
STATUS DECIMAL,
ADMISS_WARD_CODE STRING,
ADMISS_DEPT_CODE STRING,
CURR_WARD_CODE STRING,
CURR_DEPT_CODE STRING,
CURR_BED_CODE STRING,
DISCHARGE_DEPT_CODE STRING,
DISCHARGE_WARD_CODE STRING,
CHIEF_DOCTOR_CODE STRING,
ATTENDING_DOCTOR_CODE STRING,
INP_DOCTOR_CODE STRING,
TEAM_CODE STRING,
INSURANCE_NATURE_CODE STRING,
INSURANCE_NATURE_NAME STRING,
INSURANCE_PLACE_CODE STRING,
INSURANCE_PLACE_NAME STRING,
INSURANCE_TYPE_CODE STRING,
INSURANCE_TYPE_NAME STRING,
MEDICAL_DIVISION_NAME STRING,
IS_DRG DECIMAL,
MODIFY_DATE DATE,
IS_DELETE DECIMAL,
IS_ADD_POINT DECIMAL,
DRG_TYPE DECIMAL,
ODS_CREATE_TIME DATE,
ODS_UPDATE_TIME DATE,
PRIMARY KEY (PID) NOT ENFORCED
) WITH (
-- 'connector' = 'stream-x',
-- 'print' = 'true'
'connector' = 'oracle-x',
'url' = 'jdbc:oracle:thin:@//xxxx:1521/helowin',
'table-name' = 'FLINKUSER.INP_PATIENT_1',
'username' = 'xxxx',
'password' = 'xxxx',
'sink.all-replace' = 'true'
);

insert into sink_INP_PATIENT_22
select *
from source_INP_PATIENT;

+++++++++++++++++++++++++++++++++++++++++++++++

Oracle Table metadata Script:
+++++++++++++++++++++++++++++++++++++++++++++++

create table INP_PATIENT
(
HOS_CODE VARCHAR2(60),
PID VARCHAR2(128),
PKID VARCHAR2(20),
PATIENT_ID VARCHAR2(128),
HC_CARDNO VARCHAR2(128),
VISIT_CARDNO VARCHAR2(128),
INPNO VARCHAR2(128),
CASEID VARCHAR2(128),
NAME VARCHAR2(200),
BIRTHDAY DATE,
AGE NUMBER,
AGE_DESC VARCHAR2(40),
SEX_NAME VARCHAR2(40),
IDCARD_NO VARCHAR2(128),
ADMISS_NUM NUMBER,
INDAYS NUMBER,
ADMISSION_DATE DATE,
DISCHARGE_DATE DATE,
STATUS NUMBER,
ADMISS_WARD_CODE VARCHAR2(400),
ADMISS_DEPT_CODE VARCHAR2(400),
CURR_WARD_CODE VARCHAR2(400),
CURR_DEPT_CODE VARCHAR2(400),
CURR_BED_CODE VARCHAR2(400),
DISCHARGE_DEPT_CODE VARCHAR2(400),
DISCHARGE_WARD_CODE VARCHAR2(400),
CHIEF_DOCTOR_CODE VARCHAR2(400),
ATTENDING_DOCTOR_CODE VARCHAR2(400),
INP_DOCTOR_CODE VARCHAR2(400),
TEAM_CODE VARCHAR2(400),
INSURANCE_NATURE_CODE VARCHAR2(128),
INSURANCE_NATURE_NAME VARCHAR2(400),
INSURANCE_PLACE_CODE VARCHAR2(128),
INSURANCE_PLACE_NAME VARCHAR2(400),
INSURANCE_TYPE_CODE VARCHAR2(128),
INSURANCE_TYPE_NAME VARCHAR2(400),
MEDICAL_DIVISION_NAME VARCHAR2(128),
IS_DRG NUMBER,
MODIFY_DATE DATE,
IS_DELETE NUMBER,
IS_ADD_POINT NUMBER,
DRG_TYPE NUMBER,
ODS_CREATE_TIME DATE,
ODS_UPDATE_TIME DATE
)

create table INP_PATIENT_1
(
HOS_CODE VARCHAR2(60),
PID VARCHAR2(128),
PKID VARCHAR2(20),
PATIENT_ID VARCHAR2(128),
HC_CARDNO VARCHAR2(128),
VISIT_CARDNO VARCHAR2(128),
INPNO VARCHAR2(128),
CASEID VARCHAR2(128),
NAME VARCHAR2(200),
BIRTHDAY DATE,
AGE NUMBER,
AGE_DESC VARCHAR2(40),
SEX_NAME VARCHAR2(40),
IDCARD_NO VARCHAR2(128),
ADMISS_NUM NUMBER,
INDAYS NUMBER,
ADMISSION_DATE DATE,
DISCHARGE_DATE DATE,
STATUS NUMBER,
ADMISS_WARD_CODE VARCHAR2(400),
ADMISS_DEPT_CODE VARCHAR2(400),
CURR_WARD_CODE VARCHAR2(400),
CURR_DEPT_CODE VARCHAR2(400),
CURR_BED_CODE VARCHAR2(400),
DISCHARGE_DEPT_CODE VARCHAR2(400),
DISCHARGE_WARD_CODE VARCHAR2(400),
CHIEF_DOCTOR_CODE VARCHAR2(400),
ATTENDING_DOCTOR_CODE VARCHAR2(400),
INP_DOCTOR_CODE VARCHAR2(400),
TEAM_CODE VARCHAR2(400),
INSURANCE_NATURE_CODE VARCHAR2(128),
INSURANCE_NATURE_NAME VARCHAR2(400),
INSURANCE_PLACE_CODE VARCHAR2(128),
INSURANCE_PLACE_NAME VARCHAR2(400),
INSURANCE_TYPE_CODE VARCHAR2(128),
INSURANCE_TYPE_NAME VARCHAR2(400),
MEDICAL_DIVISION_NAME VARCHAR2(128),
IS_DRG NUMBER,
MODIFY_DATE DATE,
IS_DELETE NUMBER,
IS_ADD_POINT NUMBER,
DRG_TYPE NUMBER,
ODS_CREATE_TIME DATE,
ODS_UPDATE_TIME DATE
)

INSERT INTO FLINKUSER.INP_PATIENT (HOS_CODE, PID, PKID, PATIENT_ID, HC_CARDNO, VISIT_CARDNO, INPNO, CASEID, NAME, BIRTHDAY, AGE, AGE_DESC, SEX_NAME, IDCARD_NO, ADMISS_NUM, INDAYS, ADMISSION_DATE, DISCHARGE_DATE, STATUS, ADMISS_WARD_CODE, ADMISS_DEPT_CODE, CURR_WARD_CODE, CURR_DEPT_CODE, CURR_BED_CODE, DISCHARGE_DEPT_CODE, DISCHARGE_WARD_CODE, CHIEF_DOCTOR_CODE, ATTENDING_DOCTOR_CODE, INP_DOCTOR_CODE, TEAM_CODE, INSURANCE_NATURE_CODE, INSURANCE_NATURE_NAME, INSURANCE_PLACE_CODE, INSURANCE_PLACE_NAME, INSURANCE_TYPE_CODE, INSURANCE_TYPE_NAME, MEDICAL_DIVISION_NAME, IS_DRG, MODIFY_DATE, IS_DELETE, IS_ADD_POINT, DRG_TYPE, ODS_CREATE_TIME, ODS_UPDATE_TIME) VALUES ('12', 'asdwqe', 'asd', 'asd', 'asdqw', 'qweqw', 'q1', '1', '123', TO_DATE('2022-04-21 19:27:06', 'YYYY-MM-DD HH24:MI:SS'), 21, 'asd', 'ad', '12', 2, 123, TO_DATE('2022-04-06 19:27:19', 'YYYY-MM-DD HH24:MI:SS'), TO_DATE('2022-04-06 19:27:22', 'YYYY-MM-DD HH24:MI:SS'), 1, '1', '1123', '4324', '123', '124', '1', '1', '1', '1', null, '12', '1', '1123', '123', null, null, '咋', '哈哈', null, TO_DATE('2022-04-06 19:28:14', 'YYYY-MM-DD HH24:MI:SS'), null, null, null, TO_DATE('2022-04-06 19:28:25', 'YYYY-MM-DD HH24:MI:SS'), TO_DATE('2022-04-06 19:28:25', 'YYYY-MM-DD HH24:MI:SS'));

+++++++++++++++++++++++++++++++++++++++++++++++

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.