DTStack / DTStack/chunjun

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

Aperta
#679 0 commenti 0 reazioni 0 assegnatari Vedi su GitHub
bug
Lingua principale
Java
Stelle
4.1k
Fork
1.7k
Metriche di merge delle PR
Nessuna PR unita negli ultimi 30g

Descrizione

**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'));

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

Guida per i contributori

Apri la guida per i contributori

Direzione di ricerca

Inizia riproducendo il Chunjun SQL fornito con la source oraclelogminer-x e il sink oracle-x, verificando se l’evento della source viene emesso e se l’operazione del sink viene completata. Confronta le definizioni delle tabelle source e target e verifica che la riga prevista compaia in INP_PATIENT_1 senza errori.

Scritto dal modello di indicizzazione a partire dal testo della issue.

Valutazione

Stack tecnologico
java, sql
Ambito
databases
Tipo di issue
Bug
Difficoltà
4/5
Tempo stimato
3-5 giorni
Stato di attività
Ferma
Chiarezza
Da chiarire
Idoneità per principianti
35/100

Ricevi le nuove issue nella tua casella

Un breve riepilogo di issue GitHub adatte ai principianti.