@@ -52,18 +52,30 @@ def get_output_table():
5252
5353 res = get_output_table ()
5454 assert len (res ) == 3
55- assert res [0 ] == ("message1" ,)
56- assert res [1 ] == ("message2" ,)
57- assert res [2 ] == ("message3" ,)
55+ if dest_uri .startswith ("cratedb://" ):
56+ messages_db = [res [0 ][0 ], res [1 ][0 ], res [2 ][0 ]]
57+ assert "message1" in messages_db
58+ assert "message2" in messages_db
59+ assert "message3" in messages_db
60+ else :
61+ assert res [0 ] == ("message1" ,)
62+ assert res [1 ] == ("message2" ,)
63+ assert res [2 ] == ("message3" ,)
5864
5965 # run again, nothing should be inserted into the output table
6066 run ()
6167
6268 res = get_output_table ()
6369 assert len (res ) == 3
64- assert res [0 ] == ("message1" ,)
65- assert res [1 ] == ("message2" ,)
66- assert res [2 ] == ("message3" ,)
70+ if dest_uri .startswith ("cratedb://" ):
71+ messages_db = [res [0 ][0 ], res [1 ][0 ], res [2 ][0 ]]
72+ assert "message1" in messages_db
73+ assert "message2" in messages_db
74+ assert "message3" in messages_db
75+ else :
76+ assert res [0 ] == ("message1" ,)
77+ assert res [1 ] == ("message2" ,)
78+ assert res [2 ] == ("message3" ,)
6779
6880 # add a new message
6981 producer .produce (topic , "message4" .encode ("utf-8" ))
@@ -73,10 +85,17 @@ def get_output_table():
7385 run ()
7486 res = get_output_table ()
7587 assert len (res ) == 4
76- assert res [0 ] == ("message1" ,)
77- assert res [1 ] == ("message2" ,)
78- assert res [2 ] == ("message3" ,)
79- assert res [3 ] == ("message4" ,)
88+ if dest_uri .startswith ("cratedb://" ):
89+ messages_db = [res [0 ][0 ], res [1 ][0 ], res [2 ][0 ], res [3 ][0 ]]
90+ assert "message1" in messages_db
91+ assert "message2" in messages_db
92+ assert "message3" in messages_db
93+ assert "message4" in messages_db
94+ else :
95+ assert res [0 ] == ("message1" ,)
96+ assert res [1 ] == ("message2" ,)
97+ assert res [2 ] == ("message3" ,)
98+ assert res [3 ] == ("message4" ,)
8099
81100
82101@pytest .mark .parametrize (
0 commit comments