1.10版本 无法使用WATERMARK ,报空指针。
Dieses Issue hat noch niemand übernommen.
Bewertung
- Schwierigkeit
- 4/5
- Geschätzter Aufwand
- 3-5 Tage
- Anfängerfreundlichkeit
- 25/100
- Issue-Typ
- Bug
- Klarheit
- Größtenteils klar
- Aktivitätsstatus
- Veraltet
- Bereich
- stream-processing
Rechercherichtung
Beginnen Sie mit der CREATE TABLE-Anweisung des Issues und dem CustomerWaterMarkerForLong-Fehler aus dem Flink 1.10-Log. Vergleichen Sie das von UNIX_TIMESTAMP abgeleitete bigint-Feld und die WATERMARK-Syntax mit dem unterstützten Verhalten in 1.10. Die Aufgabe ist erledigt, wenn das gemeldete Schema ein vollständiges, funktionierendes WATERMARK-Beispiel enthält und den Null-Pointer-Fehler nicht mehr erzeugt.
Vom Indexierungsmodell aus dem Issue-Text verfasst.
Beschreibung
日志样例
{"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 语法案例。实际测试中无法使用
- Vorherrschende Sprache
- Java
- Sterne
- 2k
- Forks
- 913
- PR-Merge-Kennzahlen
- Keine gemergten PRs in 30 T.
Beitragsleitfaden
Für dieses Repository ist kein Beitragsleitfaden indexiert
Erste Schritte
- Lesen Sie das ganze Issue und danach den Beitragsleitfaden des Projekts.
- Schreiben Sie ins Issue, dass Sie es übernehmen — das erspart doppelte Arbeit.
- Forken Sie das Repository und arbeiten Sie in einem Branch.
- Öffnen Sie einen Pull Request, der die Issue-Nummer nennt.
Mehr aus DTStack/flinkStreamSQL
-
Schwierigkeit 5/5 Über eine Woche Anfängerfreundlichkeit 20/100
DTStack/flinkStreamSQL#467 · 1 Kommentar ·
-
Schwierigkeit 4/5 3-5 Tage Anfängerfreundlichkeit 25/100
DTStack/flinkStreamSQL#437 · 3 Kommentare ·
-
Schwierigkeit 3/5 1-2 Tage Anfängerfreundlichkeit 38/100
DTStack/flinkStreamSQL#432 · 1 Kommentar ·
-
Schwierigkeit 5/5 Über eine Woche Anfängerfreundlichkeit 15/100
DTStack/flinkStreamSQL#431 ·
-
有无对sql 语法校验的方法 Offen
Schwierigkeit 5/5 Über eine Woche Anfängerfreundlichkeit 15/100
DTStack/flinkStreamSQL#430 · 1 Kommentar ·
Alle Issues in DTStack/flinkStreamSQL
Ähnliche Issues
-
certification
Schwierigkeit 1/5 Unter einer Stunde Anfängerfreundlichkeit 80/100
-
Schwierigkeit 2/5 1-3 Stunden Anfängerfreundlichkeit 75/100
-
[BUG] ECR GetAuthorizationToken returns a proxyEndpoint for the default region, not the request's Offenbug ecr
Schwierigkeit 2/5 1-3 Stunden Anfängerfreundlichkeit 75/100
-
Needs: Triage Type: Feature request
Schwierigkeit 2/5 1-3 Stunden Anfängerfreundlichkeit 70/100
AntennaPod/AntennaPod#8794 ·
-
agentic-workflows
Schwierigkeit 2/5 1-3 Stunden Anfängerfreundlichkeit 65/100
github/copilot-sdk#2760 ·