1.10版本 无法使用WATERMARK ,报空指针。

Đang mở
#334 2 bình luận 0 reaction 0 người được giao Xem trên GitHub

Chưa có ai nhận issue này.

Đánh giá

Độ khó
4/5
Thời gian dự kiến
3-5 ngày
Mức phù hợp với người mới
25/100
Loại issue
Lỗi
Độ rõ ràng
Khá rõ ràng
Mức độ hoạt động
Đình trệ
Công nghệ
java, sql
Lĩnh vực
stream-processing

Hướng nghiên cứu

Bắt đầu với câu lệnh CREATE TABLE của issue và lỗi CustomerWaterMarkerForLong trong log Flink 1.10. So sánh trường bigint được suy ra từ UNIX_TIMESTAMP và cú pháp WATERMARK với hành vi được hỗ trợ trong 1.10. Được xem là hoàn tất khi schema được báo cáo có một ví dụ WATERMARK hoàn chỉnh, hoạt động và không còn tạo ra lỗi con trỏ null.

Do mô hình lập chỉ mục viết ra từ nội dung của issue.

Mô tả

日志样例
{"time":"2020-06-17 23:00:01.211","pushid":"pushback_send-BC110-12392212311","app":"mm"}
由于时间序列是bigint类型,用UNIX_TIMESTAMP进行转换
Flink运行日志报错:
ERROR com.dtstack.flink.sql.watermarker.CustomerWaterMarkerForLong -
java.lang.NullPointerException

建表语句:
CREATE TABLE MyTable (
time varchar ,
pushid varchar ,
app varchar ,
UNIX_TIMESTAMP(time, 'yyyy-MM-dd HH:mm:ss')*1000 bigint AS xctime ,
WATERMARK FOR xctime AS withOffset( xctime , 1000)
)
WITH (
type='kafka11',
bootstrapServers='kafka:9092',
offsetReset='latest',
topic='test_1',
groupId='flink_sql',
parallelism='4',
timezone='Asia/Shanghai',
topicIsPattern ='false',
sourcedatatype ='dt_nest'
);

CREATE TABLE MyResult(
app VARCHAR,
cnt BIGINT,
wStart timestamp
)WITH(
type ='elasticsearch6',
address ='eshost:9200',
cluster='bigdata-es5.6',
estype ='date',
index ='MyResult',
parallelism ='1',
id='0'
);

insert into MyResult
select
d.app ,
count(d.pushid) cnt ,
TUMBLE_START(d.ROWTIME, INTERVAL '3' SECOND) as wStart
from
MyTable as d
group by d.app,TUMBLE(d.ROWTIME, INTERVAL '3' SECOND);

可否提供一份完整的1.10版本 WATERMARK 语法案例。实际测试中无法使用

Ngôn ngữ chính
Java
Star
2k
Fork
913
Chỉ số merge pull request
Không có pull request nào được merge trong 30 ngày

Hướng dẫn đóng góp

Chưa lập chỉ mục được hướng dẫn đóng góp cho kho mã nguồn này

Bắt đầu từ đâu

  1. Đọc hết issue, rồi đọc hướng dẫn đóng góp của dự án.
  2. Bình luận trên issue rằng bạn sẽ nhận — tránh hai người làm cùng một việc.
  3. Fork repository và làm thay đổi trên một nhánh.
  4. Mở pull request có tham chiếu số hiệu của issue.

Issue khác của DTStack/flinkStreamSQL

Tất cả issue của DTStack/flinkStreamSQL

Issue tương tự

Thêm issue về Java

Nhận issue mới trong hộp thư của bạn

Bản tóm tắt ngắn những issue GitHub phù hợp với người mới.